diff --git a/src/iceberg/catalog/rest/CMakeLists.txt b/src/iceberg/catalog/rest/CMakeLists.txt index f64860ff4..0736290bc 100644 --- a/src/iceberg/catalog/rest/CMakeLists.txt +++ b/src/iceberg/catalog/rest/CMakeLists.txt @@ -34,6 +34,8 @@ set(ICEBERG_REST_SOURCES rest_catalog.cc rest_file_io.cc rest_metrics_reporter.cc + rest_table.cc + rest_table_scan.cc rest_util.cc types.cc) diff --git a/src/iceberg/catalog/rest/catalog_properties.cc b/src/iceberg/catalog/rest/catalog_properties.cc index 0e417e6c3..91d72e8a4 100644 --- a/src/iceberg/catalog/rest/catalog_properties.cc +++ b/src/iceberg/catalog/rest/catalog_properties.cc @@ -20,6 +20,7 @@ #include "iceberg/catalog/rest/catalog_properties.h" #include +#include #include #include @@ -61,4 +62,14 @@ Result RestCatalogProperties::SnapshotLoadingMode() const { } } +Result> RestCatalogProperties::ScanPlanningModeFrom( + const std::unordered_map& config) { + auto it = config.find(kScanPlanningMode.key()); + if (it == config.end()) return std::nullopt; + std::string lower = StringUtils::ToLower(it->second); + if (lower == "client") return ScanPlanningMode::kClient; + if (lower == "server") return ScanPlanningMode::kServer; + return InvalidArgument("Invalid scan planning mode: '{}'.", it->second); +} + } // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/catalog_properties.h b/src/iceberg/catalog/rest/catalog_properties.h index d1ee0e9c4..63dcd9a95 100644 --- a/src/iceberg/catalog/rest/catalog_properties.h +++ b/src/iceberg/catalog/rest/catalog_properties.h @@ -34,6 +34,9 @@ namespace iceberg::rest { /// \brief Snapshot loading mode for REST catalog. enum class SnapshotMode : uint8_t { kAll, kRefs }; +/// \brief Scan planning mode for REST catalog. +enum class ScanPlanningMode : uint8_t { kClient, kServer }; + /// \brief Configuration class for a REST Catalog. class ICEBERG_REST_EXPORT RestCatalogProperties : public ConfigBase { @@ -58,6 +61,8 @@ class ICEBERG_REST_EXPORT RestCatalogProperties /// \brief Whether to report metrics to the REST catalog server (default: true). inline static Entry kMetricsReportingEnabled{ "rest-metrics-reporting-enabled", "true"}; + /// \brief The scan planning mode (client or server). + inline static Entry kScanPlanningMode{"scan-planning-mode", "client"}; /// \brief The prefix for HTTP headers. inline static constexpr std::string_view kHeaderPrefix = "header."; @@ -80,6 +85,11 @@ class ICEBERG_REST_EXPORT RestCatalogProperties /// "REFS", or an error if the value is invalid. Parsing is /// case-insensitive to match Java behavior. Result SnapshotLoadingMode() const; + + /// \brief Get the scan planning mode from the given config map, returning + /// std::nullopt if the key is absent. + static Result> ScanPlanningModeFrom( + const std::unordered_map& config); }; } // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/http_client.cc b/src/iceberg/catalog/rest/http_client.cc index 6661c5098..ec870b87d 100644 --- a/src/iceberg/catalog/rest/http_client.cc +++ b/src/iceberg/catalog/rest/http_client.cc @@ -62,6 +62,15 @@ std::unordered_map HttpResponse::headers() const { return impl_->headers(); } +HttpResponse HttpResponse::MakeForTesting(int32_t status_code, std::string body) { + cpr::Response cpr_response; + cpr_response.status_code = status_code; + cpr_response.text = std::move(body); + HttpResponse response; + response.impl_ = std::make_unique(std::move(cpr_response)); + return response; +} + namespace { /// \brief Default error type for unparseable REST responses. diff --git a/src/iceberg/catalog/rest/http_client.h b/src/iceberg/catalog/rest/http_client.h index ea9c10a39..d42798fae 100644 --- a/src/iceberg/catalog/rest/http_client.h +++ b/src/iceberg/catalog/rest/http_client.h @@ -61,6 +61,9 @@ class ICEBERG_REST_EXPORT HttpResponse { /// \brief Get the headers of the response as a map. std::unordered_map headers() const; + /// \brief Create a response for use in unit tests. + static HttpResponse MakeForTesting(int32_t status_code, std::string body); + private: friend class HttpClient; class Impl; @@ -71,7 +74,7 @@ class ICEBERG_REST_EXPORT HttpResponse { class ICEBERG_REST_EXPORT HttpClient { public: explicit HttpClient(std::unordered_map default_headers = {}); - ~HttpClient(); + virtual ~HttpClient(); HttpClient(const HttpClient&) = delete; HttpClient& operator=(const HttpClient&) = delete; @@ -79,36 +82,35 @@ class ICEBERG_REST_EXPORT HttpClient { HttpClient& operator=(HttpClient&&) = delete; /// \brief Sends a GET request. - Result Get(const std::string& path, - const std::unordered_map& params, - const std::unordered_map& headers, - const ErrorHandler& error_handler, auth::AuthSession& session); + virtual Result Get( + const std::string& path, const std::unordered_map& params, + const std::unordered_map& headers, + const ErrorHandler& error_handler, auth::AuthSession& session); /// \brief Sends a POST request. - Result Post(const std::string& path, const std::string& body, - const std::unordered_map& headers, - const ErrorHandler& error_handler, - auth::AuthSession& session); + virtual Result Post( + const std::string& path, const std::string& body, + const std::unordered_map& headers, + const ErrorHandler& error_handler, auth::AuthSession& session); /// \brief Sends a POST request with form data. - Result PostForm( + virtual Result PostForm( const std::string& path, const std::unordered_map& form_data, const std::unordered_map& headers, const ErrorHandler& error_handler, auth::AuthSession& session); /// \brief Sends a HEAD request. - Result Head(const std::string& path, - const std::unordered_map& headers, - const ErrorHandler& error_handler, - auth::AuthSession& session); + virtual Result Head( + const std::string& path, + const std::unordered_map& headers, + const ErrorHandler& error_handler, auth::AuthSession& session); /// \brief Sends a DELETE request. - Result Delete(const std::string& path, - const std::unordered_map& params, - const std::unordered_map& headers, - const ErrorHandler& error_handler, - auth::AuthSession& session); + virtual Result Delete( + const std::string& path, const std::unordered_map& params, + const std::unordered_map& headers, + const ErrorHandler& error_handler, auth::AuthSession& session); private: std::unordered_map default_headers_; diff --git a/src/iceberg/catalog/rest/json_serde.cc b/src/iceberg/catalog/rest/json_serde.cc index 3ce753f18..d6520878a 100644 --- a/src/iceberg/catalog/rest/json_serde.cc +++ b/src/iceberg/catalog/rest/json_serde.cc @@ -531,6 +531,15 @@ Result ScanTaskFieldsToJson( json[kFileScanTasks] = std::move(tasks_json); } + if (!response.storage_credentials.empty()) { + nlohmann::json creds_json = nlohmann::json::array(); + for (const auto& cred : response.storage_credentials) { + ICEBERG_ASSIGN_OR_RAISE(auto entry, StorageCredentialToJson(cred)); + creds_json.push_back(std::move(entry)); + } + json[kStorageCredentials] = std::move(creds_json); + } + return json; } @@ -571,6 +580,21 @@ Status ScanTaskFieldsFromJson( FileScanTasksFromJson(file_scan_tasks_json, response.delete_files, partition_specs_by_id, schema)); } + + // 4. storage_credentials + if (json.contains(kStorageCredentials)) { + ICEBERG_ASSIGN_OR_RAISE(auto creds_json, + GetJsonValue(json, kStorageCredentials)); + if (!creds_json.is_array()) { + return JsonParseError("Cannot parse storage credentials from non-array: {}", + SafeDumpJson(creds_json)); + } + for (const auto& entry : creds_json) { + ICEBERG_ASSIGN_OR_RAISE(auto cred, StorageCredentialFromJson(entry)); + response.storage_credentials.push_back(std::move(cred)); + } + } + return {}; } diff --git a/src/iceberg/catalog/rest/meson.build b/src/iceberg/catalog/rest/meson.build index 65a67ffb8..27fb7ac34 100644 --- a/src/iceberg/catalog/rest/meson.build +++ b/src/iceberg/catalog/rest/meson.build @@ -31,6 +31,8 @@ iceberg_rest_sources = files( 'rest_catalog.cc', 'rest_file_io.cc', 'rest_metrics_reporter.cc', + 'rest_table.cc', + 'rest_table_scan.cc', 'rest_util.cc', 'types.cc', ) @@ -95,6 +97,8 @@ install_headers( 'resource_paths.h', 'rest_catalog.h', 'rest_file_io.h', + 'rest_table.h', + 'rest_table_scan.h', 'rest_util.h', 'type_fwd.h', 'types.h', diff --git a/src/iceberg/catalog/rest/rest_catalog.cc b/src/iceberg/catalog/rest/rest_catalog.cc index 4a4f990ea..5a8bb2507 100644 --- a/src/iceberg/catalog/rest/rest_catalog.cc +++ b/src/iceberg/catalog/rest/rest_catalog.cc @@ -39,9 +39,12 @@ #include "iceberg/catalog/rest/resource_paths.h" #include "iceberg/catalog/rest/rest_file_io.h" #include "iceberg/catalog/rest/rest_metrics_reporter_internal.h" +#include "iceberg/catalog/rest/rest_table.h" #include "iceberg/catalog/rest/rest_util.h" #include "iceberg/catalog/rest/types.h" #include "iceberg/json_serde_internal.h" +#include "iceberg/logging/log_level.h" +#include "iceberg/logging/logger.h" #include "iceberg/metrics/metrics_reporters.h" #include "iceberg/partition_spec.h" #include "iceberg/result.h" @@ -454,7 +457,7 @@ Result> RestCatalog::Make( RestCatalog::RestCatalog(RestCatalogProperties config, std::shared_ptr file_io, std::shared_ptr client, - std::unique_ptr paths, + std::shared_ptr paths, std::unordered_set endpoints, std::unique_ptr auth_manager, std::shared_ptr catalog_session, @@ -899,6 +902,47 @@ Result> RestCatalog::MakeTableFromLoadResult( auto table_catalog = std::make_shared( shared_from_this(), context, identifier, table_config, table_session, table_io); + // Determine effective scan planning mode: table config overrides client config. + ICEBERG_ASSIGN_OR_RAISE(auto client_mode, + RestCatalogProperties::ScanPlanningModeFrom(config_.configs())); + ICEBERG_ASSIGN_OR_RAISE(auto server_mode, + RestCatalogProperties::ScanPlanningModeFrom(table_config)); + + if (client_mode.has_value() && server_mode.has_value() && + *client_mode != *server_mode) { + Log(LogLevel::kWarn, + "Scan planning mode mismatch for table {}: client config={}, server config={}. " + "Server config will take precedence.", + identifier.ToString(), + *client_mode == ScanPlanningMode::kClient ? "client" : "server", + *server_mode == ScanPlanningMode::kClient ? "client" : "server"); + } + + ScanPlanningMode effective_mode = + server_mode.value_or(client_mode.value_or(ScanPlanningMode::kClient)); + + if (effective_mode == ScanPlanningMode::kServer) { + if (!supported_endpoints_.contains(Endpoint::PlanTableScan())) { + return NotSupported( + "Server requires server-side scan planning for table {} but does not support " + "the PlanTableScan endpoint.", + identifier.ToString()); + } + RestScanContext rest_ctx{ + .client = client_, + .paths = paths_, + .session = table_session, + .supported_endpoints = supported_endpoints_, + .identifier = identifier, + .catalog_config = config_.configs(), + .table_config = table_config, + }; + return RestTable::Make(identifier, std::move(result.metadata), + std::move(result.metadata_location), std::move(table_io), + std::move(table_catalog), RestTableName(name_, identifier), + reporter, std::move(rest_ctx)); + } + return Table::Make(identifier, std::move(result.metadata), std::move(result.metadata_location), std::move(table_io), std::move(table_catalog), RestTableName(name_, identifier), diff --git a/src/iceberg/catalog/rest/rest_catalog.h b/src/iceberg/catalog/rest/rest_catalog.h index 65b0b5eab..8b194b668 100644 --- a/src/iceberg/catalog/rest/rest_catalog.h +++ b/src/iceberg/catalog/rest/rest_catalog.h @@ -69,7 +69,7 @@ class ICEBERG_REST_EXPORT RestCatalog final class TableScopedCatalog; RestCatalog(RestCatalogProperties config, std::shared_ptr file_io, - std::shared_ptr client, std::unique_ptr paths, + std::shared_ptr client, std::shared_ptr paths, std::unordered_set endpoints, std::unique_ptr auth_manager, std::shared_ptr catalog_session, @@ -193,7 +193,7 @@ class ICEBERG_REST_EXPORT RestCatalog final RestCatalogProperties config_; std::shared_ptr file_io_; std::shared_ptr client_; - std::unique_ptr paths_; + std::shared_ptr paths_; std::string name_; std::unordered_set supported_endpoints_; std::unique_ptr auth_manager_; diff --git a/src/iceberg/catalog/rest/rest_table.cc b/src/iceberg/catalog/rest/rest_table.cc new file mode 100644 index 000000000..6fb330670 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table.cc @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/catalog/rest/rest_table.h" + +#include +#include + +#include "iceberg/catalog/rest/rest_table_scan.h" +#include "iceberg/result.h" +#include "iceberg/table_metadata.h" +#include "iceberg/util/macros.h" + +namespace iceberg::rest { + +RestTable::RestTable(TableIdentifier identifier, std::shared_ptr metadata, + std::string metadata_location, std::shared_ptr io, + std::shared_ptr catalog, std::string full_name, + std::shared_ptr reporter, + RestScanContext rest_context) + : Table(std::move(identifier), std::move(metadata), std::move(metadata_location), + std::move(io), std::move(catalog), std::move(full_name), std::move(reporter)), + rest_context_(std::move(rest_context)) {} + +RestTable::~RestTable() = default; + +Result> RestTable::Make( + TableIdentifier identifier, std::shared_ptr metadata, + std::string metadata_location, std::shared_ptr io, + std::shared_ptr catalog, std::string full_name, + std::shared_ptr reporter, RestScanContext rest_context) { + if (metadata == nullptr) { + return InvalidArgument("Metadata cannot be null"); + } + return std::shared_ptr( + new RestTable(std::move(identifier), std::move(metadata), + std::move(metadata_location), std::move(io), std::move(catalog), + std::move(full_name), std::move(reporter), std::move(rest_context))); +} + +Result> RestTable::NewScan() const { + return std::make_unique(metadata_, io_, full_name_, reporter_, + rest_context_); +} + +} // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/rest_table.h b/src/iceberg/catalog/rest/rest_table.h new file mode 100644 index 000000000..fb23ca66b --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table.h @@ -0,0 +1,62 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +#include +#include + +#include "iceberg/catalog/rest/iceberg_rest_export.h" +#include "iceberg/catalog/rest/rest_table_scan.h" +#include "iceberg/metrics/metrics_reporter.h" +#include "iceberg/result.h" +#include "iceberg/table.h" +#include "iceberg/type_fwd.h" + +/// \file iceberg/catalog/rest/rest_table.h +/// A Table subclass that uses server-side distributed scan planning via the REST catalog. + +namespace iceberg::rest { + +/// \brief A Table whose NewScan() returns a RestTableScanBuilder, delegating +/// PlanFiles() to the REST catalog server's scan planning endpoints. +class ICEBERG_REST_EXPORT RestTable final : public Table { + public: + static Result> Make( + TableIdentifier identifier, std::shared_ptr metadata, + std::string metadata_location, std::shared_ptr io, + std::shared_ptr catalog, std::string full_name, + std::shared_ptr reporter, RestScanContext rest_context); + + ~RestTable() override; + + /// \brief Returns a RestTableScanBuilder that will delegate PlanFiles() to the + /// REST catalog server. + Result> NewScan() const override; + + private: + RestTable(TableIdentifier identifier, std::shared_ptr metadata, + std::string metadata_location, std::shared_ptr io, + std::shared_ptr catalog, std::string full_name, + std::shared_ptr reporter, RestScanContext rest_context); + + RestScanContext rest_context_; +}; + +} // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/rest_table_scan.cc b/src/iceberg/catalog/rest/rest_table_scan.cc new file mode 100644 index 000000000..3e11c71dc --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table_scan.cc @@ -0,0 +1,294 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/catalog/rest/rest_table_scan.h" + +#include +#include + +#include + +#include "iceberg/catalog/rest/endpoint.h" +#include "iceberg/catalog/rest/error_handlers.h" +#include "iceberg/catalog/rest/http_client.h" +#include "iceberg/catalog/rest/json_serde_internal.h" +#include "iceberg/catalog/rest/resource_paths.h" +#include "iceberg/catalog/rest/rest_file_io.h" +#include "iceberg/catalog/rest/types.h" +#include "iceberg/json_serde_internal.h" +#include "iceberg/partition_spec.h" +#include "iceberg/result.h" +#include "iceberg/schema.h" +#include "iceberg/table_metadata.h" +#include "iceberg/util/macros.h" + +namespace iceberg::rest { + +namespace { + +constexpr int64_t kMinSleepMs = 1'000; +constexpr int64_t kMaxSleepMs = 60'000; +constexpr int kMaxRetries = 10; +constexpr int64_t kMaxWaitTimeMs = 5 * 60 * 1'000; + +#define ICEBERG_ENDPOINT_CHECK(endpoints, endpoint) \ + do { \ + if (!endpoints.contains(endpoint)) { \ + return NotSupported("Not supported endpoint: {}", endpoint.ToString()); \ + } \ + } while (0) + +} // namespace + +// RestTableScan + +RestTableScan::RestTableScan(std::shared_ptr metadata, + std::shared_ptr schema, std::shared_ptr io, + internal::TableScanContext context, + RestScanContext rest_context) + : DataTableScan(std::move(metadata), std::move(schema), std::move(io), + std::move(context)), + rest_context_(std::move(rest_context)) {} + +Result> RestTableScan::Make( + std::shared_ptr metadata, std::shared_ptr schema, + std::shared_ptr io, internal::TableScanContext context, + RestScanContext rest_context) { + ICEBERG_PRECHECK(metadata != nullptr, "Table metadata cannot be null"); + ICEBERG_PRECHECK(schema != nullptr, "Schema cannot be null"); + ICEBERG_PRECHECK(io != nullptr, "FileIO cannot be null"); + return std::unique_ptr( + new RestTableScan(std::move(metadata), std::move(schema), std::move(io), + std::move(context), std::move(rest_context))); +} + +Result>> RestTableScan::PlanFiles() const { + TableMetadataCache metadata_cache(metadata_.get()); + ICEBERG_ASSIGN_OR_RAISE(auto specs_by_id, metadata_cache.GetPartitionSpecsById()); + + std::string plan_id; + return PlanTableScan(plan_id, specs_by_id); +} + +Result>> RestTableScan::PlanTableScan( + std::string& plan_id, + const std::unordered_map>& specs) const { + ICEBERG_ENDPOINT_CHECK(rest_context_.supported_endpoints, Endpoint::PlanTableScan()); + + // Build request from scan context + PlanTableScanRequest request; + request.select = context_.selected_columns; + request.filter = context_.filter; + request.case_sensitive = context_.case_sensitive; + request.min_rows_requested = context_.min_rows_requested; + + if (context_.from_snapshot_id.has_value() && context_.to_snapshot_id.has_value()) { + request.start_snapshot_id = context_.from_snapshot_id; + request.end_snapshot_id = context_.to_snapshot_id; + request.use_snapshot_schema = true; + } else if (context_.snapshot_id.has_value()) { + request.snapshot_id = context_.snapshot_id; + request.use_snapshot_schema = context_.use_snapshot_schema; + } + + if (!context_.columns_to_keep_stats.empty()) { + for (int32_t field_id : context_.columns_to_keep_stats) { + ICEBERG_ASSIGN_OR_RAISE(auto name, schema_->FindColumnNameById(field_id)); + if (name.has_value()) { + request.stats_fields.emplace_back(*name); + } + } + } + + ICEBERG_ASSIGN_OR_RAISE(auto path, rest_context_.paths->Plan(rest_context_.identifier)); + ICEBERG_ASSIGN_OR_RAISE(auto request_json, ToJson(request)); + ICEBERG_ASSIGN_OR_RAISE(auto json_request, ToJsonString(request_json)); + ICEBERG_ASSIGN_OR_RAISE( + const auto response, + rest_context_.client->Post(path, json_request, /*headers=*/{}, + *PlanErrorHandler::Instance(), *rest_context_.session)); + ICEBERG_ASSIGN_OR_RAISE(auto json, FromJsonString(response.body())); + ICEBERG_ASSIGN_OR_RAISE(auto result, + PlanTableScanResponseFromJson(json, specs, *schema_)); + ICEBERG_RETURN_UNEXPECTED(result.Validate()); + + plan_id = result.plan_id; + + switch (result.plan_status) { + case PlanStatus::kCompleted: { + ICEBERG_RETURN_UNEXPECTED(ApplyStorageCredentials(result.storage_credentials)); + auto tasks = ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs); + if (!tasks.has_value()) CancelPlanning(plan_id); + return tasks; + } + case PlanStatus::kSubmitted: + return FetchPlanningResult(plan_id, specs); + case PlanStatus::kFailed: + return IOError("Scan planning failed: {}", + result.error ? result.error->message : "unknown error"); + case PlanStatus::kCancelled: + return IOError("Scan planning was cancelled for plan_id={}", plan_id); + } + return IOError("Unexpected plan status"); +} + +Result>> RestTableScan::FetchPlanningResult( + const std::string& plan_id, + const std::unordered_map>& specs) const { + ICEBERG_ENDPOINT_CHECK(rest_context_.supported_endpoints, + Endpoint::FetchPlanningResult()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, + rest_context_.paths->Plan(rest_context_.identifier, plan_id)); + + auto delay_ms = kMinSleepMs; + auto start = std::chrono::steady_clock::now(); + + for (int retry = 0; retry <= kMaxRetries; ++retry) { + ICEBERG_ASSIGN_OR_RAISE( + const auto response, + rest_context_.client->Get(path, /*params=*/{}, /*headers=*/{}, + *PlanErrorHandler::Instance(), *rest_context_.session)); + ICEBERG_ASSIGN_OR_RAISE(auto json, FromJsonString(response.body())); + ICEBERG_ASSIGN_OR_RAISE(auto result, + FetchPlanningResultResponseFromJson(json, specs, *schema_)); + ICEBERG_RETURN_UNEXPECTED(result.Validate()); + + switch (result.plan_status) { + case PlanStatus::kCompleted: { + ICEBERG_RETURN_UNEXPECTED(ApplyStorageCredentials(result.storage_credentials)); + auto tasks = ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs); + if (!tasks.has_value()) CancelPlanning(plan_id); + return tasks; + } + case PlanStatus::kSubmitted: { + auto elapsed_ms = std::chrono::duration_cast( + std::chrono::steady_clock::now() - start) + .count(); + if (elapsed_ms >= kMaxWaitTimeMs) { + CancelPlanning(plan_id); + return IOError("Scan planning timed out after {}ms waiting for plan_id={}", + elapsed_ms, plan_id); + } + std::this_thread::sleep_for(std::chrono::milliseconds(delay_ms)); + delay_ms = std::min(delay_ms * 2, kMaxSleepMs); + continue; + } + case PlanStatus::kFailed: + CancelPlanning(plan_id); + return IOError("Scan planning failed: {}", + result.error ? result.error->message : "unknown error"); + case PlanStatus::kCancelled: + return IOError("Scan planning was cancelled for plan_id={}", plan_id); + } + } + + CancelPlanning(plan_id); + return IOError("Scan planning exceeded max retries ({}) for plan_id={}", kMaxRetries, + plan_id); +} + +Result>> RestTableScan::FetchScanTasks( + const std::string& plan_task, + const std::unordered_map>& specs) const { + ICEBERG_ENDPOINT_CHECK(rest_context_.supported_endpoints, Endpoint::FetchScanTasks()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, + rest_context_.paths->FetchScanTasks(rest_context_.identifier)); + FetchScanTasksRequest request{.planTask = plan_task}; + ICEBERG_ASSIGN_OR_RAISE(auto json_request, ToJsonString(ToJson(request))); + ICEBERG_ASSIGN_OR_RAISE(const auto response, + rest_context_.client->Post(path, json_request, /*headers=*/{}, + *PlanTaskErrorHandler::Instance(), + *rest_context_.session)); + ICEBERG_ASSIGN_OR_RAISE(auto json, FromJsonString(response.body())); + ICEBERG_ASSIGN_OR_RAISE(auto result, + FetchScanTasksResponseFromJson(json, specs, *schema_)); + ICEBERG_RETURN_UNEXPECTED(result.Validate()); + ICEBERG_RETURN_UNEXPECTED(ApplyStorageCredentials(result.storage_credentials)); + + return ResolveScanTasks(result.plan_tasks, result.file_scan_tasks, specs); +} + +Result>> RestTableScan::ResolveScanTasks( + const std::optional>& plan_tasks, + const std::optional>>& file_scan_tasks, + const std::unordered_map>& specs) const { + std::vector> result; + + if (file_scan_tasks.has_value()) { + result.insert(result.end(), file_scan_tasks->begin(), file_scan_tasks->end()); + } + + if (plan_tasks.has_value()) { + for (const auto& plan_task : *plan_tasks) { + ICEBERG_ASSIGN_OR_RAISE(auto tasks, FetchScanTasks(plan_task, specs)); + result.insert(result.end(), tasks.begin(), tasks.end()); + } + } + + return result; +} + +void RestTableScan::CancelPlanning(const std::string& plan_id) const { + if (plan_id.empty()) return; + if (!rest_context_.supported_endpoints.contains(Endpoint::CancelPlanning())) return; + + auto path = rest_context_.paths->Plan(rest_context_.identifier, plan_id); + if (!path.has_value()) return; + + // Best-effort: ignore errors. + std::ignore = + rest_context_.client->Delete(*path, /*params=*/{}, /*headers=*/{}, + *PlanErrorHandler::Instance(), *rest_context_.session); +} + +const std::shared_ptr& RestTableScan::effective_io() const { + return scan_io_ ? scan_io_ : io_; +} + +Status RestTableScan::ApplyStorageCredentials( + const std::vector& credentials) const { + if (credentials.empty()) return {}; + ICEBERG_ASSIGN_OR_RAISE( + auto io, MakeTableFileIO(rest_context_.catalog_config, rest_context_.table_config, + credentials)); + scan_io_ = std::move(io); + return {}; +} + +// RestTableScanBuilder + +RestTableScanBuilder::RestTableScanBuilder( + std::shared_ptr metadata, std::shared_ptr io, + std::string table_name, std::shared_ptr metrics_reporter, + RestScanContext rest_context) + : DataTableScanBuilder(std::move(metadata), std::move(io), std::move(table_name), + std::move(metrics_reporter)), + rest_context_(std::move(rest_context)) {} + +Result> RestTableScanBuilder::Build() { + ICEBERG_RETURN_UNEXPECTED(CheckErrors()); + ICEBERG_RETURN_UNEXPECTED(context_.Validate()); + ICEBERG_ASSIGN_OR_RAISE(auto schema, ResolveSnapshotSchema()); + return RestTableScan::Make(metadata_, schema.get(), io_, std::move(context_), + rest_context_); +} + +} // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/rest_table_scan.h b/src/iceberg/catalog/rest/rest_table_scan.h new file mode 100644 index 000000000..63a852303 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table_scan.h @@ -0,0 +1,137 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#pragma once + +#include +#include +#include +#include +#include + +#include "iceberg/catalog/rest/endpoint.h" +#include "iceberg/catalog/rest/iceberg_rest_export.h" +#include "iceberg/metrics/metrics_reporter.h" +#include "iceberg/result.h" +#include "iceberg/storage_credential.h" +#include "iceberg/table_identifier.h" +#include "iceberg/table_scan.h" +#include "iceberg/type_fwd.h" + +/// \file iceberg/catalog/rest/rest_table_scan.h +/// REST-specific table scan that delegates scan planning to the REST catalog server. + +namespace iceberg::rest { + +class HttpClient; +class ResourcePaths; + +namespace auth { +class AuthSession; +} // namespace auth + +/// \brief HTTP context shared between RestTable and RestTableScan. +struct ICEBERG_REST_EXPORT RestScanContext { + std::shared_ptr client; + std::shared_ptr paths; + std::shared_ptr session; + std::unordered_set supported_endpoints; + TableIdentifier identifier; + /// Catalog-level config, used with table_config to build a scan-scoped FileIO + /// when the server vends storage credentials in a planning response. + std::unordered_map catalog_config; + /// Table-level config merged with catalog_config for scan-scoped FileIO creation. + std::unordered_map table_config; +}; + +/// \brief A DataTableScan that delegates PlanFiles() to the REST catalog server +/// via the scan planning endpoints (planTableScan / fetchPlanningResult / +/// cancelPlanning / fetchScanTasks). +class ICEBERG_REST_EXPORT RestTableScan : public DataTableScan { + public: + ~RestTableScan() override = default; + + static Result> Make( + std::shared_ptr metadata, std::shared_ptr schema, + std::shared_ptr io, internal::TableScanContext context, + RestScanContext rest_context); + + /// \brief Plans files via the REST scan planning endpoints. + Result>> PlanFiles() const override; + + /// \brief Returns the FileIO to use when reading scan results. + /// + /// If the server vended storage credentials during planning, returns a FileIO + /// initialised with those credentials; otherwise returns the table's FileIO. + /// Must be called after PlanFiles(). + const std::shared_ptr& effective_io() const; + + private: + RestTableScan(std::shared_ptr metadata, std::shared_ptr schema, + std::shared_ptr io, internal::TableScanContext context, + RestScanContext rest_context); + + /// POST /plan → handle COMPLETED / SUBMITTED / FAILED / CANCELLED. + Result>> PlanTableScan( + std::string& plan_id, + const std::unordered_map>& specs) const; + + /// GET /plan/{plan_id} with exponential backoff until COMPLETED. + Result>> FetchPlanningResult( + const std::string& plan_id, + const std::unordered_map>& specs) const; + + /// POST /tasks/{plan_task_id} → fetch FileScanTasks for one opaque plan task token. + Result>> FetchScanTasks( + const std::string& plan_task, + const std::unordered_map>& specs) const; + + /// Flatten plan_tasks (opaque tokens) + file_scan_tasks into a single list. + Result>> ResolveScanTasks( + const std::optional>& plan_tasks, + const std::optional>>& file_scan_tasks, + const std::unordered_map>& specs) const; + + /// DELETE /plan/{plan_id}; best-effort, errors are silently ignored. + void CancelPlanning(const std::string& plan_id) const; + + /// Builds a scan-scoped FileIO from vended credentials and caches it in scan_io_. + /// No-op if credentials is empty. + Status ApplyStorageCredentials(const std::vector& credentials) const; + + RestScanContext rest_context_; + mutable std::shared_ptr scan_io_; +}; + +/// \brief Builder that produces a RestTableScan with the REST HTTP context injected. +class ICEBERG_REST_EXPORT RestTableScanBuilder : public DataTableScanBuilder { + public: + RestTableScanBuilder(std::shared_ptr metadata, + std::shared_ptr io, std::string table_name, + std::shared_ptr metrics_reporter, + RestScanContext rest_context); + + /// \brief Resolves schema/context via parent logic then creates a RestTableScan. + Result> Build() override; + + private: + RestScanContext rest_context_; +}; + +} // namespace iceberg::rest diff --git a/src/iceberg/catalog/rest/types.cc b/src/iceberg/catalog/rest/types.cc index 84fba9a7c..ef20f44da 100644 --- a/src/iceberg/catalog/rest/types.cc +++ b/src/iceberg/catalog/rest/types.cc @@ -210,6 +210,7 @@ bool OptionalSharedPtrVectorEqual( template bool ScanTaskFieldsEqual(const Response& lhs, const Response& rhs) { return lhs.plan_tasks == rhs.plan_tasks && + lhs.storage_credentials == rhs.storage_credentials && SharedPtrVectorEqual(lhs.delete_files, rhs.delete_files) && OptionalSharedPtrVectorEqual(lhs.file_scan_tasks, rhs.file_scan_tasks, FileScanTaskEqual); diff --git a/src/iceberg/catalog/rest/types.h b/src/iceberg/catalog/rest/types.h index 20a59fa59..ad58127c7 100644 --- a/src/iceberg/catalog/rest/types.h +++ b/src/iceberg/catalog/rest/types.h @@ -337,7 +337,7 @@ struct ICEBERG_REST_EXPORT PlanTableScanResponse { PlanStatus plan_status = PlanStatus::kCompleted; std::string plan_id; std::optional error; - // TODO(sandeepg): Add storage credentials and bind scan FileIO to them. + std::vector storage_credentials; Status Validate() const; @@ -352,7 +352,7 @@ struct ICEBERG_REST_EXPORT FetchPlanningResultResponse { std::vector> delete_files; PlanStatus plan_status = PlanStatus::kCompleted; std::optional error; - // TODO(sandeepg): Add storage credentials and bind scan FileIO to them. + std::vector storage_credentials; Status Validate() const; @@ -373,6 +373,7 @@ struct ICEBERG_REST_EXPORT FetchScanTasksResponse { std::optional> plan_tasks; std::optional>> file_scan_tasks; std::vector> delete_files; + std::vector storage_credentials; Status Validate() const; diff --git a/src/iceberg/table_scan.cc b/src/iceberg/table_scan.cc index abc487a6d..6aa70858f 100644 --- a/src/iceberg/table_scan.cc +++ b/src/iceberg/table_scan.cc @@ -317,6 +317,7 @@ TableScanBuilder& TableScanBuilder::UseSnapshot(int64_t snap context_.snapshot_id.value()); ICEBERG_BUILDER_ASSIGN_OR_RETURN(std::ignore, metadata_->SnapshotById(snapshot_id)); context_.snapshot_id = snapshot_id; + context_.use_snapshot_schema = true; return *this; } @@ -324,6 +325,7 @@ template TableScanBuilder& TableScanBuilder::UseRef(const std::string& ref) { if (ref == SnapshotRef::kMainBranch) { context_.snapshot_id.reset(); + context_.use_snapshot_schema = false; return *this; } @@ -336,6 +338,7 @@ TableScanBuilder& TableScanBuilder::UseRef(const std::string const int64_t snapshot_id = iter->second->snapshot_id; ICEBERG_BUILDER_ASSIGN_OR_RETURN(std::ignore, metadata_->SnapshotById(snapshot_id)); context_.snapshot_id = snapshot_id; + context_.use_snapshot_schema = (iter->second->type() == SnapshotRefType::kTag); return *this; } diff --git a/src/iceberg/table_scan.h b/src/iceberg/table_scan.h index bee2b7d1d..ce374bf1a 100644 --- a/src/iceberg/table_scan.h +++ b/src/iceberg/table_scan.h @@ -212,7 +212,7 @@ class ICEBERG_EXPORT DeletedDataFileScanTask : public ChangelogScanTask { namespace internal { // Internal table scan context used by different scan implementations. -struct TableScanContext { +struct ICEBERG_EXPORT TableScanContext { std::optional snapshot_id; std::shared_ptr filter; bool ignore_residuals{false}; @@ -227,6 +227,7 @@ struct TableScanContext { std::optional to_snapshot_id; std::string branch{}; std::optional min_rows_requested; + bool use_snapshot_schema{false}; OptionalExecutor plan_executor; std::string table_name; std::shared_ptr metrics_reporter; @@ -389,13 +390,16 @@ class ICEBERG_TEMPLATE_CLASS_EXPORT TableScanBuilder : public ErrorCollector { /// \brief Builds and returns a TableScan instance. /// \return A Result containing the TableScan or an error. - Result> Build(); + virtual Result> Build(); protected: TableScanBuilder(std::shared_ptr metadata, std::shared_ptr io, std::string table_name, std::shared_ptr metrics_reporter); + TableScanBuilder(TableScanBuilder&&) = default; + TableScanBuilder& operator=(TableScanBuilder&&) = default; + // Return the schema bound to the specified snapshot. Result>> ResolveSnapshotSchema(); Status ResolveColumnStatsSelection(); @@ -461,7 +465,7 @@ class ICEBERG_EXPORT DataTableScan : public TableScan { /// \brief Plans the scan tasks by resolving manifests and data files. /// \return A Result containing scan tasks or an error. - Result>> PlanFiles() const; + virtual Result>> PlanFiles() const; private: Status ReportScan(const Snapshot& snapshot, const ScanMetrics& scan_metrics) const; diff --git a/src/iceberg/test/CMakeLists.txt b/src/iceberg/test/CMakeLists.txt index 5ca9fd915..e21d75c59 100644 --- a/src/iceberg/test/CMakeLists.txt +++ b/src/iceberg/test/CMakeLists.txt @@ -313,11 +313,13 @@ if(ICEBERG_BUILD_REST) add_rest_iceberg_test(rest_catalog_test SOURCES auth_manager_test.cc + catalog_properties_test.cc error_handlers_test.cc endpoint_test.cc rest_file_io_test.cc rest_json_serde_test.cc rest_metrics_reporter_test.cc + rest_table_scan_test.cc rest_util_test.cc) if(ICEBERG_SIGV4) diff --git a/src/iceberg/test/catalog_properties_test.cc b/src/iceberg/test/catalog_properties_test.cc new file mode 100644 index 000000000..5f22e60d4 --- /dev/null +++ b/src/iceberg/test/catalog_properties_test.cc @@ -0,0 +1,81 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/catalog/rest/catalog_properties.h" + +#include + +#include "iceberg/test/matchers.h" + +namespace iceberg::rest { + +TEST(ScanPlanningModeTest, MissingKeyReturnsNullopt) { + std::unordered_map config; + auto result = RestCatalogProperties::ScanPlanningModeFrom(config); + ASSERT_THAT(result, IsOk()); + EXPECT_FALSE(result->has_value()); +} + +TEST(ScanPlanningModeTest, ClientLowercaseReturnsKClient) { + auto result = + RestCatalogProperties::ScanPlanningModeFrom({{"scan-planning-mode", "client"}}); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result->has_value()); + EXPECT_EQ(**result, ScanPlanningMode::kClient); +} + +TEST(ScanPlanningModeTest, ServerLowercaseReturnsKServer) { + auto result = + RestCatalogProperties::ScanPlanningModeFrom({{"scan-planning-mode", "server"}}); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result->has_value()); + EXPECT_EQ(**result, ScanPlanningMode::kServer); +} + +TEST(ScanPlanningModeTest, ClientUppercaseReturnsKClient) { + auto result = + RestCatalogProperties::ScanPlanningModeFrom({{"scan-planning-mode", "CLIENT"}}); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result->has_value()); + EXPECT_EQ(**result, ScanPlanningMode::kClient); +} + +TEST(ScanPlanningModeTest, ServerUppercaseReturnsKServer) { + auto result = + RestCatalogProperties::ScanPlanningModeFrom({{"scan-planning-mode", "SERVER"}}); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result->has_value()); + EXPECT_EQ(**result, ScanPlanningMode::kServer); +} + +TEST(ScanPlanningModeTest, InvalidValueReturnsError) { + auto result = + RestCatalogProperties::ScanPlanningModeFrom({{"scan-planning-mode", "invalid"}}); + EXPECT_THAT(result, IsError(ErrorKind::kInvalidArgument)); +} + +TEST(ScanPlanningModeTest, OtherKeysAreIgnored) { + auto result = RestCatalogProperties::ScanPlanningModeFrom( + {{"other-key", "server"}, {"scan-planning-mode", "client"}}); + ASSERT_THAT(result, IsOk()); + ASSERT_TRUE(result->has_value()); + EXPECT_EQ(**result, ScanPlanningMode::kClient); +} + +} // namespace iceberg::rest diff --git a/src/iceberg/test/meson.build b/src/iceberg/test/meson.build index 6844df5e9..4773941ec 100644 --- a/src/iceberg/test/meson.build +++ b/src/iceberg/test/meson.build @@ -147,11 +147,13 @@ if get_option('rest').enabled() 'rest_catalog_test': { 'sources': files( 'auth_manager_test.cc', + 'catalog_properties_test.cc', 'endpoint_test.cc', 'error_handlers_test.cc', 'rest_file_io_test.cc', 'rest_json_serde_test.cc', 'rest_metrics_reporter_test.cc', + 'rest_table_scan_test.cc', 'rest_util_test.cc', ), 'dependencies': [iceberg_rest_dep], diff --git a/src/iceberg/test/rest_table_scan_test.cc b/src/iceberg/test/rest_table_scan_test.cc new file mode 100644 index 000000000..9c5efd929 --- /dev/null +++ b/src/iceberg/test/rest_table_scan_test.cc @@ -0,0 +1,467 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +#include "iceberg/catalog/rest/rest_table_scan.h" + +#include +#include +#include +#include +#include + +#include +#include +#include + +#include "iceberg/catalog/rest/auth/auth_session.h" +#include "iceberg/catalog/rest/endpoint.h" +#include "iceberg/catalog/rest/error_handlers.h" +#include "iceberg/catalog/rest/http_client.h" +#include "iceberg/catalog/rest/resource_paths.h" +#include "iceberg/catalog/rest/rest_table.h" +#include "iceberg/file_io.h" +#include "iceberg/partition_spec.h" +#include "iceberg/schema.h" +#include "iceberg/snapshot.h" +#include "iceberg/table_identifier.h" +#include "iceberg/table_metadata.h" +#include "iceberg/table_scan.h" +#include "iceberg/test/matchers.h" +#include "iceberg/type.h" + +namespace iceberg::rest { + +using ::testing::_; +using ::testing::Return; + +// -------------------------------------------------------------------------- +// Mock HTTP client that overrides the virtual methods of HttpClient. +// The base class constructor creates a cpr::ConnectionPool, which is a +// lightweight allocation (no network connections are opened at construction). +// -------------------------------------------------------------------------- +class MockHttpClient : public HttpClient { + public: + MockHttpClient() : HttpClient({}) {} + + MOCK_METHOD(Result, Get, + (const std::string& path, + (const std::unordered_map&)params, + (const std::unordered_map&)headers, + const ErrorHandler& error_handler, auth::AuthSession& session), + (override)); + + MOCK_METHOD(Result, Post, + (const std::string& path, const std::string& body, + (const std::unordered_map&)headers, + const ErrorHandler& error_handler, auth::AuthSession& session), + (override)); + + MOCK_METHOD(Result, Delete, + (const std::string& path, + (const std::unordered_map&)params, + (const std::unordered_map&)headers, + const ErrorHandler& error_handler, auth::AuthSession& session), + (override)); +}; + +// -------------------------------------------------------------------------- +// Minimal FileIO stub (no real I/O needed for server-side scan planning tests) +// -------------------------------------------------------------------------- +class NoOpFileIO : public FileIO { + public: + Result ReadFile(const std::string&, std::optional) override { + return IOError("NoOpFileIO"); + } + Status WriteFile(const std::string&, std::string_view) override { return {}; } + Status DeleteFile(const std::string&) override { return {}; } +}; + +// -------------------------------------------------------------------------- +// Test fixture shared by RestTableScan tests. +// -------------------------------------------------------------------------- +class RestTableScanTest : public ::testing::Test { + protected: + void SetUp() override { + schema_ = std::make_shared( + std::vector{SchemaField::MakeRequired(1, "id", int32()), + SchemaField::MakeRequired(2, "data", string())}); + + auto spec = PartitionSpec::Unpartitioned(); + + constexpr int64_t kSnapshotId = 1000L; + auto snapshot = std::make_shared( + Snapshot{.snapshot_id = kSnapshotId, + .sequence_number = 1L, + .timestamp_ms = TimePointMsFromUnixMs(1609459200000L), + .manifest_list = "/tmp/manifest-list.avro", + .schema_id = schema_->schema_id()}); + + metadata_ = std::make_shared( + TableMetadata{.format_version = 2, + .table_uuid = "test-uuid", + .location = "/tmp/table", + .last_sequence_number = 1L, + .last_updated_ms = TimePointMsFromUnixMs(1609459200000L), + .last_column_id = 2, + .schemas = {schema_}, + .current_schema_id = schema_->schema_id(), + .partition_specs = {spec}, + .default_spec_id = spec->spec_id(), + .last_partition_id = 999, + .current_snapshot_id = kSnapshotId, + .snapshots = {snapshot}, + .refs = {{"main", std::make_shared(SnapshotRef{ + .snapshot_id = kSnapshotId, + .retention = SnapshotRef::Branch{}})}}}); + + file_io_ = std::make_shared(); + + mock_client_ = std::make_shared(); + + ICEBERG_UNWRAP_OR_FAIL(paths_, + ResourcePaths::Make("http://test-server", /*prefix=*/"", + /*namespace_separator=*/"%1F")); + + session_ = auth::AuthSession::MakeDefault(/*headers=*/{}); + + identifier_ = TableIdentifier{.ns = Namespace{{"default"}}, .name = "my_table"}; + + all_plan_endpoints_ = {Endpoint::PlanTableScan(), Endpoint::FetchPlanningResult(), + Endpoint::CancelPlanning(), Endpoint::FetchScanTasks()}; + } + + // Pass std::nullopt to get the full set of plan endpoints (default). + // Pass an explicit set (including empty) to use exactly that set. + RestScanContext MakeContext( + std::optional> endpoints = std::nullopt) { + auto effective = endpoints.has_value() ? std::move(*endpoints) : all_plan_endpoints_; + return RestScanContext{ + .client = mock_client_, + .paths = paths_, + .session = session_, + .supported_endpoints = std::move(effective), + .identifier = identifier_, + }; + } + + Result> MakeScan(RestScanContext ctx) { + return RestTableScan::Make(metadata_, schema_, file_io_, internal::TableScanContext{}, + std::move(ctx)); + } + + std::shared_ptr schema_; + std::shared_ptr metadata_; + std::shared_ptr file_io_; + std::shared_ptr mock_client_; + std::shared_ptr paths_; + std::shared_ptr session_; + TableIdentifier identifier_; + std::unordered_set all_plan_endpoints_; +}; + +// -------------------------------------------------------------------------- +// PlanFiles: server returns COMPLETED immediately, no file scan tasks. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesCompleted) { + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: server returns COMPLETED with a non-empty plan-id (still valid). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesCompletedWithPlanId) { + constexpr std::string_view kResponseBody = + R"({"status":"completed","plan-id":"plan-abc"})"; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: server returns SUBMITTED → poll returns COMPLETED. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesSubmittedThenCompleted) { + constexpr std::string_view kSubmittedBody = + R"({"status":"submitted","plan-id":"plan-poll-1"})"; + constexpr std::string_view kCompletedBody = R"({"status":"completed"})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kSubmittedBody)))); + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kCompletedBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: server returns FAILED → scan returns IOError. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesFailed) { + constexpr std::string_view kFailedBody = + R"({"status":"failed","error":{"message":"server error","type":"ServerError","code":500}})"; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kFailedBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// PlanFiles: PlanTableScan endpoint missing → NotSupported error. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesEndpointNotSupported) { + ICEBERG_UNWRAP_OR_FAIL(auto scan, + MakeScan(MakeContext(std::unordered_set{}))); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); +} + +// -------------------------------------------------------------------------- +// PlanFiles with plan-tasks: server returns COMPLETED with opaque task token, +// then FetchScanTasks is called and returns no file scan tasks. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesWithPlanTasks) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-1","plan-tasks":["tok-1"]})"; + // FetchScanTasksResponse requires at least one of plan-tasks or file-scan-tasks + // present. + constexpr std::string_view kTasksResponse = R"({"file-scan-tasks":[]})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kTasksResponse)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Cancel is called when FetchScanTasks fails after a COMPLETED response that +// included plan-tasks. This mirrors the Java cancelPlan-on-close behavior. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, CancelCalledWhenFetchScanTasksFails) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-cancel-1","plan-tasks":["tok-a"]})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(IOError("FetchScanTasks failed"))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// Cancel is a no-op when plan_id is empty (server returned no plan-id). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, CancelIsNoOpWithEmptyPlanId) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-tasks":["tok-b"]})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(IOError("FetchScanTasks failed"))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)).Times(0); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// Cancel is a no-op when CancelPlanning endpoint is not advertised. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, CancelIsNoOpWhenEndpointNotAdvertised) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-2","plan-tasks":["tok-c"]})"; + + std::unordered_set endpoints_without_cancel = { + Endpoint::PlanTableScan(), Endpoint::FetchPlanningResult(), + Endpoint::FetchScanTasks()}; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(IOError("FetchScanTasks failed"))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)).Times(0); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext(endpoints_without_cancel))); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// FetchPlanningResult: FetchPlanningResult endpoint missing → NotSupported. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, FetchPlanningResultEndpointNotSupported) { + constexpr std::string_view kSubmittedBody = + R"({"status":"submitted","plan-id":"plan-3"})"; + + std::unordered_set endpoints_without_fetch = { + Endpoint::PlanTableScan(), Endpoint::CancelPlanning(), Endpoint::FetchScanTasks()}; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kSubmittedBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext(endpoints_without_fetch))); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); +} + +// -------------------------------------------------------------------------- +// FetchScanTasks: endpoint missing → NotSupported. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, FetchScanTasksEndpointNotSupported) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-4","plan-tasks":["tok-d"]})"; + + std::unordered_set endpoints_without_tasks = {Endpoint::PlanTableScan(), + Endpoint::FetchPlanningResult(), + Endpoint::CancelPlanning()}; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext(endpoints_without_tasks))); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); +} + +// -------------------------------------------------------------------------- +// use_snapshot_schema: UseSnapshot() sets it to true in the builder context. +// RestTableScanBuilder propagates context from DataTableScanBuilder. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, UseSnapshotPropagatesUseSnapshotSchemaInContext) { + constexpr int64_t kSnapshotId = 1000L; + RestTableScanBuilder builder(metadata_, file_io_, "test.my_table", nullptr, + MakeContext(std::nullopt)); + builder.UseSnapshot(kSnapshotId); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder.Build()); + EXPECT_TRUE(scan->context().use_snapshot_schema); +} + +// -------------------------------------------------------------------------- +// use_snapshot_schema: default scan does not set use_snapshot_schema. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, DefaultScanDoesNotSetUseSnapshotSchema) { + RestTableScanBuilder builder(metadata_, file_io_, "test.my_table", nullptr, + MakeContext(std::nullopt)); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder.Build()); + EXPECT_FALSE(scan->context().use_snapshot_schema); +} + +// -------------------------------------------------------------------------- +// Storage credentials in COMPLETED response: effective_io() returns a +// credential-scoped IO, not the original table IO. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, StorageCredentialsInPlanResponseUpdatesEffectiveIO) { + constexpr std::string_view kResponseBody = R"({ + "status": "completed", + "storage-credentials": [ + {"prefix": "s3://bucket/prefix", "config": {"key": "value"}} + ] + })"; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); + + auto* rest_scan = dynamic_cast(scan.get()); + ASSERT_NE(rest_scan, nullptr); + // effective_io() must return a credential-scoped IO, not the original file_io_. + EXPECT_NE(rest_scan->effective_io().get(), file_io_.get()); +} + +// -------------------------------------------------------------------------- +// No storage credentials: effective_io() falls back to the table's FileIO. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, NoStorageCredentialsEffectiveIoFallsBackToTableIO) { + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); + + auto* rest_scan = dynamic_cast(scan.get()); + ASSERT_NE(rest_scan, nullptr); + EXPECT_EQ(rest_scan->effective_io().get(), file_io_.get()); +} + +// -------------------------------------------------------------------------- +// Storage credentials returned in FetchScanTasksResponse also update +// effective_io(). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, StorageCredentialsInFetchScanTasksResponseUpdatesEffectiveIO) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-cred","plan-tasks":["tok-cred"]})"; + constexpr std::string_view kTasksResponse = R"({ + "file-scan-tasks": [], + "storage-credentials": [ + {"prefix": "s3://bucket/prefix", "config": {"key": "value"}} + ] + })"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kTasksResponse)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); + + auto* rest_scan = dynamic_cast(scan.get()); + ASSERT_NE(rest_scan, nullptr); + EXPECT_NE(rest_scan->effective_io().get(), file_io_.get()); +} + +// -------------------------------------------------------------------------- +// RestTable::NewScan returns a RestTableScanBuilder (not a plain builder). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, RestTableNewScanReturnsRestTableScanBuilder) { + ICEBERG_UNWRAP_OR_FAIL( + auto table, RestTable::Make(identifier_, metadata_, "/tmp/metadata.json", file_io_, + /*catalog=*/nullptr, "test.my_table", nullptr, + MakeContext(std::nullopt))); + ICEBERG_UNWRAP_OR_FAIL(auto builder, table->NewScan()); + auto* typed = dynamic_cast(builder.get()); + EXPECT_NE(typed, nullptr); +} + +} // namespace iceberg::rest diff --git a/src/iceberg/test/table_scan_test.cc b/src/iceberg/test/table_scan_test.cc index 375c9aa53..bb55182eb 100644 --- a/src/iceberg/test/table_scan_test.cc +++ b/src/iceberg/test/table_scan_test.cc @@ -774,6 +774,40 @@ TEST_P(TableScanTest, SchemaWithSelectedColumnsAndFilter) { } } +// use_snapshot_schema propagation tests: verify the field is set correctly for +// UseSnapshot, UseRef (tag vs branch), and default/incremental scans. +TEST_P(TableScanTest, UseSnapshotSetsTrueUseSnapshotSchema) { + constexpr int64_t kSnapshotId = 1000L; + ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder(table_metadata_)); + builder->UseSnapshot(kSnapshotId); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build()); + EXPECT_TRUE(scan->context().use_snapshot_schema); +} + +TEST_P(TableScanTest, UseRefTagSetsTrueUseSnapshotSchema) { + constexpr int64_t kSnapshotId = 1000L; + table_metadata_->refs["v1.0"] = std::make_shared( + SnapshotRef{.snapshot_id = kSnapshotId, .retention = SnapshotRef::Tag{}}); + + ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder(table_metadata_)); + builder->UseRef("v1.0"); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build()); + EXPECT_TRUE(scan->context().use_snapshot_schema); +} + +TEST_P(TableScanTest, UseRefBranchSetsFalseUseSnapshotSchema) { + ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder(table_metadata_)); + builder->UseRef("main"); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build()); + EXPECT_FALSE(scan->context().use_snapshot_schema); +} + +TEST_P(TableScanTest, DefaultScanHasFalseUseSnapshotSchema) { + ICEBERG_UNWRAP_OR_FAIL(auto builder, MakeScanBuilder(table_metadata_)); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder->Build()); + EXPECT_FALSE(scan->context().use_snapshot_schema); +} + INSTANTIATE_TEST_SUITE_P(TableScanVersions, TableScanTest, testing::Values(1, 2, 3)); } // namespace iceberg