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/rest_catalog.cc b/src/iceberg/catalog/rest/rest_catalog.cc index 4a4f990ea..e583fbc1c 100644 --- a/src/iceberg/catalog/rest/rest_catalog.cc +++ b/src/iceberg/catalog/rest/rest_catalog.cc @@ -28,6 +28,7 @@ #include +#include "iceberg/catalog/rest/auth/auth_manager_internal.h" #include "iceberg/catalog/rest/auth/auth_managers.h" #include "iceberg/catalog/rest/auth/auth_session.h" #include "iceberg/catalog/rest/catalog_properties.h" @@ -39,9 +40,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" @@ -452,9 +456,25 @@ Result> RestCatalog::Make( snapshot_mode, std::move(default_context), std::move(reporter), metrics_executor)); } +Result> RestCatalog::MakeForTesting( + RestCatalogProperties config, std::shared_ptr file_io, + std::shared_ptr client, std::shared_ptr paths, + std::unordered_set endpoints) { + std::string catalog_name = config.Get(RestCatalogProperties::kName); + ICEBERG_ASSIGN_OR_RAISE(auto auth_manager, + auth::MakeNoopAuthManager(catalog_name, config.configs())); + auto session = auth::AuthSession::MakeDefault({}); + ICEBERG_ASSIGN_OR_RAISE(auto snapshot_mode, config.SnapshotLoadingMode()); + return std::shared_ptr( + new RestCatalog(std::move(config), std::move(file_io), std::move(client), + std::move(paths), std::move(endpoints), std::move(auth_manager), + std::move(session), snapshot_mode, SessionContext::Empty(), + /*reporter=*/nullptr, /*metrics_executor=*/nullptr)); +} + 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 +919,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..52a033178 100644 --- a/src/iceberg/catalog/rest/rest_catalog.h +++ b/src/iceberg/catalog/rest/rest_catalog.h @@ -58,6 +58,13 @@ class ICEBERG_REST_EXPORT RestCatalog final static Result> Make(const RestCatalogProperties& config, Executor* metrics_executor = nullptr); + /// \brief Test-only factory that constructs a RestCatalog with pre-built dependencies, + /// bypassing the FetchServerConfig HTTP exchange performed by Make(). + static Result> MakeForTesting( + RestCatalogProperties config, std::shared_ptr file_io, + std::shared_ptr client, std::shared_ptr paths, + std::unordered_set endpoints); + std::string_view name() const override; Result> AsCatalog() override; @@ -69,7 +76,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 +200,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..dd2b43b96 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table.cc @@ -0,0 +1,68 @@ +/* + * 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_); +} + +Result> +RestTable::NewIncrementalAppendScan() 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..f47dbe5b3 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table.h @@ -0,0 +1,66 @@ +/* + * 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; + + /// \brief Returns a RestIncrementalAppendScanBuilder that delegates to the server. + Result> NewIncrementalAppendScan() + 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..b321f4251 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table_scan.cc @@ -0,0 +1,552 @@ +/* + * 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/constants.h" +#include "iceberg/json_serde_internal.h" +#include "iceberg/partition_spec.h" +#include "iceberg/result.h" +#include "iceberg/schema.h" +#include "iceberg/snapshot.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; + +using SpecsById = std::unordered_map>; +// A shared slot that lets the scan and the stream share a single FileIO reference. +// Credentials vended by any lazy FetchScanTasks response are written into *slot, +// making them visible to RestTableScan::io() even after the stream is consumed. +using ScanIoSlot = std::shared_ptr>; + +#define ICEBERG_ENDPOINT_CHECK(endpoints, endpoint) \ + do { \ + if (!endpoints.contains(endpoint)) { \ + return NotSupported("Not supported endpoint: {}", endpoint.ToString()); \ + } \ + } while (0) + +// --------------------------------------------------------------------------- +// Shared HTTP scan planning helpers used by all REST scan implementations. +// --------------------------------------------------------------------------- + +Status ApplyStorageCredentials(const RestScanContext& ctx, + const std::vector& credentials, + std::shared_ptr& scan_io) { + if (credentials.empty()) return {}; + ICEBERG_ASSIGN_OR_RAISE( + auto io, MakeTableFileIO(ctx.catalog_config, ctx.table_config, credentials)); + scan_io = std::move(io); + return {}; +} + +void CancelPlanning(const RestScanContext& ctx, const std::string& plan_id) { + if (plan_id.empty()) return; + if (!ctx.supported_endpoints.contains(Endpoint::CancelPlanning())) return; + + auto path = ctx.paths->Plan(ctx.identifier, plan_id); + if (!path.has_value()) return; + + std::ignore = ctx.client->Delete(*path, /*params=*/{}, /*headers=*/{}, + *PlanErrorHandler::Instance(), *ctx.session); +} + +Result>> FetchScanTasks( + const RestScanContext& ctx, const Schema& schema, const std::string& plan_task, + const SpecsById& specs, std::shared_ptr& scan_io); + +Result>> ResolveScanTasks( + const RestScanContext& ctx, const Schema& schema, + const std::optional>& plan_tasks, + const std::optional>>& file_scan_tasks, + const SpecsById& specs, std::shared_ptr& scan_io) { + 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& token : *plan_tasks) { + ICEBERG_ASSIGN_OR_RAISE(auto tasks, + FetchScanTasks(ctx, schema, token, specs, scan_io)); + result.insert(result.end(), tasks.begin(), tasks.end()); + } + } + + return result; +} + +Result>> FetchScanTasks( + const RestScanContext& ctx, const Schema& schema, const std::string& plan_task, + const SpecsById& specs, std::shared_ptr& scan_io) { + ICEBERG_ENDPOINT_CHECK(ctx.supported_endpoints, Endpoint::FetchScanTasks()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, ctx.paths->FetchScanTasks(ctx.identifier)); + FetchScanTasksRequest request{.planTask = plan_task}; + ICEBERG_ASSIGN_OR_RAISE(auto json_request, ToJsonString(ToJson(request))); + ICEBERG_ASSIGN_OR_RAISE( + const auto response, + ctx.client->Post(path, json_request, /*headers=*/{}, + *PlanTaskErrorHandler::Instance(), *ctx.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(ctx, result.storage_credentials, scan_io)); + + return ResolveScanTasks(ctx, schema, result.plan_tasks, result.file_scan_tasks, specs, + scan_io); +} + +Result>> FetchPlanningResult( + const RestScanContext& ctx, const Schema& schema, const std::string& plan_id, + const SpecsById& specs, std::shared_ptr& scan_io) { + ICEBERG_ENDPOINT_CHECK(ctx.supported_endpoints, Endpoint::FetchPlanningResult()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, ctx.paths->Plan(ctx.identifier, plan_id)); + + auto delay_ms = kMinSleepMs; + auto start = std::chrono::steady_clock::now(); + + for (int retry = 0; retry <= kMaxRetries; ++retry) { + auto response_or = ctx.client->Get(path, /*params=*/{}, /*headers=*/{}, + *PlanErrorHandler::Instance(), *ctx.session); + if (!response_or) { + CancelPlanning(ctx, plan_id); + return std::unexpected(response_or.error()); + } + + auto json_or = FromJsonString(response_or->body()); + if (!json_or) { + CancelPlanning(ctx, plan_id); + return std::unexpected(json_or.error()); + } + + auto result_or = FetchPlanningResultResponseFromJson(*json_or, specs, schema); + if (!result_or) { + CancelPlanning(ctx, plan_id); + return std::unexpected(result_or.error()); + } + + if (auto s = result_or->Validate(); !s) { + CancelPlanning(ctx, plan_id); + return std::unexpected(s.error()); + } + + auto& result = *result_or; + + switch (result.plan_status) { + case PlanStatus::kCompleted: { + if (auto s = ApplyStorageCredentials(ctx, result.storage_credentials, scan_io); + !s) { + CancelPlanning(ctx, plan_id); + return std::unexpected(s.error()); + } + auto tasks = ResolveScanTasks(ctx, schema, result.plan_tasks, + result.file_scan_tasks, specs, scan_io); + if (!tasks) CancelPlanning(ctx, 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(ctx, 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(ctx, 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(ctx, plan_id); + return IOError("Scan planning exceeded max retries ({}) for plan_id={}", kMaxRetries, + plan_id); +} + +/// A lazy stream that drives FetchScanTasks calls on demand. +/// +/// The eager POST /plan (or poll until COMPLETED) is done before this stream is +/// constructed. The stream then yields the directly-returned file_scan_tasks first, +/// then fetches each plan_task token lazily one at a time as Next() is called. +/// +/// scan_io_slot is shared with the owning RestTableScan so that credentials vended +/// by any FetchScanTasks response are visible through RestTableScan::io() even after +/// the stream has been consumed. +class RestFileScanTaskStream final : public FileScanTaskStream { + public: + RestFileScanTaskStream(RestScanContext ctx, std::shared_ptr schema, + std::string plan_id, + std::vector> initial_tasks, + std::vector plan_task_tokens, SpecsById specs, + ScanIoSlot scan_io_slot) + : ctx_(std::move(ctx)), + schema_(std::move(schema)), + plan_id_(std::move(plan_id)), + buffer_(std::move(initial_tasks)), + plan_task_tokens_(std::move(plan_task_tokens)), + specs_(std::move(specs)), + scan_io_slot_(std::move(scan_io_slot)) {} + + ~RestFileScanTaskStream() override { + if (!consumed_) CancelPlanning(ctx_, plan_id_); + } + + protected: + Result>> NextImpl() override { + while (true) { + if (buffer_pos_ < buffer_.size()) { + return buffer_[buffer_pos_++]; + } + if (token_pos_ >= plan_task_tokens_.size()) { + consumed_ = true; + return std::nullopt; + } + const auto& token = plan_task_tokens_[token_pos_++]; + ICEBERG_ASSIGN_OR_RAISE( + buffer_, FetchScanTasks(ctx_, *schema_, token, specs_, *scan_io_slot_)); + buffer_pos_ = 0; + } + } + + private: + RestScanContext ctx_; + std::shared_ptr schema_; + std::string plan_id_; + std::vector> buffer_; + size_t buffer_pos_ = 0; + std::vector plan_task_tokens_; + size_t token_pos_ = 0; + SpecsById specs_; + ScanIoSlot scan_io_slot_; + bool consumed_ = false; +}; + +/// POST /plan (polling if SUBMITTED), apply credentials, return a lazy stream. +/// +/// scan_io_slot is shared with the caller (RestTableScan) so that credentials +/// vended by FetchScanTasks responses update RestTableScan::io() in place. +Result ExecuteScanPlanStream( + const RestScanContext& ctx, const Schema& schema, std::shared_ptr schema_ptr, + PlanTableScanRequest request, const SpecsById& specs, ScanIoSlot scan_io_slot) { + ICEBERG_ENDPOINT_CHECK(ctx.supported_endpoints, Endpoint::PlanTableScan()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, ctx.paths->Plan(ctx.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, + ctx.client->Post(path, json_request, /*headers=*/{}, + *PlanErrorHandler::Instance(), *ctx.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()); + + const std::string plan_id = result.plan_id; + + if (result.plan_status == PlanStatus::kSubmitted) { + // Poll until COMPLETED, eagerly collecting all tasks into the buffer. + std::string mutable_plan_id = plan_id; + ICEBERG_ASSIGN_OR_RAISE(auto tasks, FetchPlanningResult(ctx, schema, mutable_plan_id, + specs, *scan_io_slot)); + return std::make_unique( + ctx, schema_ptr, mutable_plan_id, std::move(tasks), std::vector{}, + specs, scan_io_slot); + } + + if (result.plan_status == PlanStatus::kFailed) { + CancelPlanning(ctx, plan_id); + return IOError("Scan planning failed: {}", + result.error ? result.error->message : "unknown error"); + } + if (result.plan_status == PlanStatus::kCancelled) { + return IOError("Scan planning was cancelled for plan_id={}", plan_id); + } + + // kCompleted: apply credentials from the initial response, then build the lazy stream. + if (auto s = ApplyStorageCredentials(ctx, result.storage_credentials, *scan_io_slot); + !s) { + CancelPlanning(ctx, plan_id); + return std::unexpected(s.error()); + } + + std::vector> initial_tasks; + if (result.file_scan_tasks.has_value()) { + initial_tasks = std::move(*result.file_scan_tasks); + } + std::vector plan_task_tokens; + if (result.plan_tasks.has_value()) { + plan_task_tokens = std::move(*result.plan_tasks); + } + + return std::make_unique( + ctx, std::move(schema_ptr), plan_id, std::move(initial_tasks), + std::move(plan_task_tokens), specs, std::move(scan_io_slot)); +} + +/// Eager batch planning: used by RestIncrementalAppendScan which has no stream path. +Result>> ExecuteScanPlan( + const RestScanContext& ctx, const Schema& schema, PlanTableScanRequest request, + const SpecsById& specs, std::shared_ptr& scan_io) { + ICEBERG_ENDPOINT_CHECK(ctx.supported_endpoints, Endpoint::PlanTableScan()); + + ICEBERG_ASSIGN_OR_RAISE(auto path, ctx.paths->Plan(ctx.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, + ctx.client->Post(path, json_request, /*headers=*/{}, + *PlanErrorHandler::Instance(), *ctx.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()); + + const std::string plan_id = result.plan_id; + + switch (result.plan_status) { + case PlanStatus::kCompleted: { + if (auto s = ApplyStorageCredentials(ctx, result.storage_credentials, scan_io); + !s) { + CancelPlanning(ctx, plan_id); + return std::unexpected(s.error()); + } + auto tasks = ResolveScanTasks(ctx, schema, result.plan_tasks, + result.file_scan_tasks, specs, scan_io); + if (!tasks) CancelPlanning(ctx, plan_id); + return tasks; + } + case PlanStatus::kSubmitted: + return FetchPlanningResult(ctx, schema, plan_id, specs, scan_io); + case PlanStatus::kFailed: + CancelPlanning(ctx, 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); + } + return IOError("Unexpected plan status"); +} + +} // 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)), + scan_io_slot_(std::make_shared>()) {} + +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::PlanFilesStream() const { + *scan_io_slot_ = + nullptr; // reset so stale credentials from a prior plan are not reused + TableMetadataCache metadata_cache(metadata_.get()); + ICEBERG_ASSIGN_OR_RAISE(auto specs, metadata_cache.GetPartitionSpecsById()); + + PlanTableScanRequest request; + request.select = context_.selected_columns.value_or(std::vector{}); + 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); + } + } + } + + return ExecuteScanPlanStream(rest_context_, *schema_, schema_, std::move(request), + specs, scan_io_slot_); +} + +const std::shared_ptr& RestTableScan::io() const { + return *scan_io_slot_ ? *scan_io_slot_ : io_; +} + +// --------------------------------------------------------------------------- +// 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_); +} + +// --------------------------------------------------------------------------- +// RestIncrementalAppendScan +// --------------------------------------------------------------------------- + +RestIncrementalAppendScan::RestIncrementalAppendScan( + std::shared_ptr metadata, std::shared_ptr schema, + std::shared_ptr io, internal::TableScanContext context, + RestScanContext rest_context) + : IncrementalAppendScan(std::move(metadata), std::move(schema), std::move(io), + std::move(context)), + rest_context_(std::move(rest_context)) {} + +Result> RestIncrementalAppendScan::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 RestIncrementalAppendScan(std::move(metadata), std::move(schema), std::move(io), + std::move(context), std::move(rest_context))); +} + +Result>> RestIncrementalAppendScan::PlanFiles() + const { + scan_io_ = nullptr; // reset so stale credentials from a prior plan are not reused + TableMetadataCache metadata_cache(metadata_.get()); + ICEBERG_ASSIGN_OR_RAISE(auto specs, metadata_cache.GetPartitionSpecsById()); + + PlanTableScanRequest request; + request.select = context_.selected_columns.value_or(std::vector{}); + request.filter = context_.filter; + request.case_sensitive = context_.case_sensitive; + request.min_rows_requested = context_.min_rows_requested; + + // Resolve end snapshot: use to_snapshot_id if set, else current table snapshot. + if (context_.to_snapshot_id.has_value()) { + request.end_snapshot_id = context_.to_snapshot_id; + } else { + if (metadata_->current_snapshot_id == kInvalidSnapshotId) return {}; + ICEBERG_ASSIGN_OR_RAISE(auto snapshot, metadata_->Snapshot()); + if (!snapshot) return {}; + request.end_snapshot_id = snapshot->snapshot_id; + } + + // Resolve start snapshot (exclusive): respect from_snapshot_id_inclusive. + if (context_.from_snapshot_id.has_value()) { + if (context_.from_snapshot_id_inclusive) { + ICEBERG_ASSIGN_OR_RAISE(auto from_snap, + metadata_->SnapshotById(*context_.from_snapshot_id)); + request.start_snapshot_id = from_snap->parent_snapshot_id; + } else { + request.start_snapshot_id = context_.from_snapshot_id; + } + } + + 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); + } + } + } + + return ExecuteScanPlan(rest_context_, *schema_, std::move(request), specs, scan_io_); +} + +// --------------------------------------------------------------------------- +// RestIncrementalAppendScanBuilder +// --------------------------------------------------------------------------- + +RestIncrementalAppendScanBuilder::RestIncrementalAppendScanBuilder( + std::shared_ptr metadata, std::shared_ptr io, + std::string table_name, std::shared_ptr metrics_reporter, + RestScanContext rest_context) + : IncrementalAppendScanBuilder(std::move(metadata), std::move(io), + std::move(table_name), std::move(metrics_reporter)), + rest_context_(std::move(rest_context)) {} + +Result> RestIncrementalAppendScanBuilder::Build() { + ICEBERG_RETURN_UNEXPECTED(CheckErrors()); + ICEBERG_RETURN_UNEXPECTED(context_.Validate()); + ICEBERG_ASSIGN_OR_RAISE(auto schema, ResolveSnapshotSchema()); + return RestIncrementalAppendScan::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..8df976959 --- /dev/null +++ b/src/iceberg/catalog/rest/rest_table_scan.h @@ -0,0 +1,153 @@ +/* + * 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 scans that delegate 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 REST scan classes. +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 PlanFilesStream() 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 lazily via the REST scan planning endpoints. + Result PlanFilesStream() const override; + + /// \brief Returns the effective FileIO for reading scan results. + /// + /// If the server vended storage credentials during planning, returns a FileIO + /// initialised with those credentials; otherwise returns the table's FileIO. + const std::shared_ptr& io() const override; + + private: + RestTableScan(std::shared_ptr metadata, std::shared_ptr schema, + std::shared_ptr io, internal::TableScanContext context, + RestScanContext rest_context); + + RestScanContext rest_context_; + /// Shared slot so credentials vended by any lazy FetchScanTasks response are + /// visible through io() even after the stream has been consumed. + mutable std::shared_ptr> scan_io_slot_; +}; + +/// \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_; +}; + +/// \brief An IncrementalAppendScan that delegates PlanFiles() to the REST catalog server. +/// +/// IncrementalChangelogScan is not delegated because the REST planTableScan response +/// only carries FileScanTask objects; reconstructing ChangelogScanTask entries +/// (AddedRowsScanTask vs DeletedDataFileScanTask) requires per-snapshot operation +/// metadata that the server does not return. Changelog scans always plan locally. +class ICEBERG_REST_EXPORT RestIncrementalAppendScan : public IncrementalAppendScan { + public: + ~RestIncrementalAppendScan() 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; + + private: + RestIncrementalAppendScan(std::shared_ptr metadata, + std::shared_ptr schema, std::shared_ptr io, + internal::TableScanContext context, + RestScanContext rest_context); + + RestScanContext rest_context_; + mutable std::shared_ptr scan_io_; +}; + +/// \brief Builder that produces a RestIncrementalAppendScan. +class ICEBERG_REST_EXPORT RestIncrementalAppendScanBuilder + : public IncrementalAppendScanBuilder { + public: + RestIncrementalAppendScanBuilder(std::shared_ptr metadata, + std::shared_ptr io, std::string table_name, + std::shared_ptr metrics_reporter, + RestScanContext rest_context); + + 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..7384e2e39 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); @@ -297,10 +298,10 @@ Status PlanTableScanResponse::Validate() const { "Invalid response: tasks can only be defined when status is 'completed'"); } if (!plan_id.empty() && plan_status != PlanStatus::kSubmitted && - plan_status != PlanStatus::kCompleted) { + plan_status != PlanStatus::kCompleted && plan_status != PlanStatus::kFailed) { return ValidationFailed( - "Invalid response: plan id can only be defined when status is 'submitted' or " - "'completed'"); + "Invalid response: plan id can only be defined when status is 'submitted', " + "'completed', or 'failed'"); } if (!HasNonEmptyFileScanTasks(*this) && !delete_files.empty()) { return ValidationFailed( 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 1c68d5341..7ce3b4d87 100644 --- a/src/iceberg/table_scan.cc +++ b/src/iceberg/table_scan.cc @@ -408,6 +408,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; } @@ -415,6 +416,7 @@ template TableScanBuilder& TableScanBuilder::UseRef(const std::string& ref) { if (ref == SnapshotRef::kMainBranch) { context_.snapshot_id.reset(); + context_.use_snapshot_schema = false; return *this; } @@ -427,6 +429,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 7310d435b..0190ad01f 100644 --- a/src/iceberg/table_scan.h +++ b/src/iceberg/table_scan.h @@ -219,7 +219,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}; @@ -234,6 +234,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; @@ -402,13 +403,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(); @@ -438,7 +442,7 @@ class ICEBERG_EXPORT TableScan { const internal::TableScanContext& context() const; /// \brief Returns the file I/O instance used for reading files. - const std::shared_ptr& io() const; + virtual const std::shared_ptr& io() const; /// \brief Returns this scan's filter expression. const std::shared_ptr& filter() const; @@ -484,7 +488,7 @@ class ICEBERG_EXPORT DataTableScan : public TableScan { /// can outlive this scan. An executor configured through PlanWith() is borrowed and /// must remain alive until the stream is destroyed, as later Next() calls may submit /// work to it. - Result PlanFilesStream() const; + virtual Result PlanFilesStream() const; protected: using TableScan::TableScan; diff --git a/src/iceberg/test/CMakeLists.txt b/src/iceberg/test/CMakeLists.txt index f23bc9181..70694187a 100644 --- a/src/iceberg/test/CMakeLists.txt +++ b/src/iceberg/test/CMakeLists.txt @@ -318,11 +318,14 @@ 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_catalog_unit_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/rest_catalog_unit_test.cc b/src/iceberg/test/rest_catalog_unit_test.cc new file mode 100644 index 000000000..b23d92a77 --- /dev/null +++ b/src/iceberg/test/rest_catalog_unit_test.cc @@ -0,0 +1,211 @@ +/* + * 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 +#include +#include +#include + +#include +#include + +#include "iceberg/catalog/rest/catalog_properties.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_catalog.h" +#include "iceberg/catalog/rest/rest_table.h" +#include "iceberg/file_io.h" +#include "iceberg/table_identifier.h" +#include "iceberg/test/matchers.h" + +namespace iceberg::rest { + +using ::testing::_; +using ::testing::Return; + +// -------------------------------------------------------------------------- +// Mock HTTP client (same pattern as rest_table_scan_test.cc) +// -------------------------------------------------------------------------- +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 +// -------------------------------------------------------------------------- +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 {}; } +}; + +// -------------------------------------------------------------------------- +// Helper JSON bodies for LoadTable GET responses +// -------------------------------------------------------------------------- + +// Minimal table metadata blob (no snapshots, no partitions). +constexpr std::string_view kBaseMetadataJson = + R"("metadata-location":"s3://bucket/metadata/v1.json","metadata":{"format-version":2,"table-uuid":"test-uuid","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":{}})"; + +// LoadTable response where server config requests server-side scan planning. +const std::string kLoadTableServerScanResponse = + std::string(R"({"config":{"scan-planning-mode":"server"},)") + + std::string(kBaseMetadataJson) + "}"; + +// LoadTable response with no scan-planning-mode set (defaults to client). +const std::string kLoadTableDefaultResponse = + std::string("{") + std::string(kBaseMetadataJson) + "}"; + +// -------------------------------------------------------------------------- +// Test fixture: builds a RestCatalog via MakeForTesting so no HTTP calls +// are made during catalog construction. +// -------------------------------------------------------------------------- +class RestCatalogLoadTableTest : public ::testing::Test { + protected: + void SetUp() override { + mock_client_ = std::make_shared(); + file_io_ = std::make_shared(); + + ICEBERG_UNWRAP_OR_FAIL(paths_, + ResourcePaths::Make("http://test-server", /*prefix=*/"", + /*namespace_separator=*/"%1F")); + + identifier_ = TableIdentifier{.ns = Namespace{{"default"}}, .name = "my_table"}; + + all_plan_endpoints_ = {Endpoint::LoadTable(), Endpoint::PlanTableScan(), + Endpoint::FetchPlanningResult(), Endpoint::CancelPlanning(), + Endpoint::FetchScanTasks()}; + + no_plan_endpoint_set_ = {Endpoint::LoadTable()}; + } + + // Builds a RestCatalog with the given supported endpoints and an optional + // client-side scan-planning-mode setting. + Result> MakeCatalog( + const std::unordered_set& endpoints, + const std::string& client_scan_mode = "") { + std::unordered_map props{ + {"uri", "http://test-server"}, + }; + if (!client_scan_mode.empty()) { + props["scan-planning-mode"] = client_scan_mode; + } + auto config = RestCatalogProperties::FromMap(props); + return RestCatalog::MakeForTesting(std::move(config), file_io_, mock_client_, paths_, + endpoints); + } + + std::shared_ptr mock_client_; + std::shared_ptr file_io_; + std::shared_ptr paths_; + TableIdentifier identifier_; + std::unordered_set all_plan_endpoints_; + std::unordered_set no_plan_endpoint_set_; +}; + +// -------------------------------------------------------------------------- +// Table config "scan-planning-mode":"server" → LoadTable returns a RestTable. +// -------------------------------------------------------------------------- +TEST_F(RestCatalogLoadTableTest, TableConfigServerScanReturnsRestTable) { + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, kLoadTableServerScanResponse))); + + ICEBERG_UNWRAP_OR_FAIL(auto catalog, MakeCatalog(all_plan_endpoints_)); + ICEBERG_UNWRAP_OR_FAIL(auto as_catalog, catalog->AsCatalog()); + ICEBERG_UNWRAP_OR_FAIL(auto table, as_catalog->LoadTable(identifier_)); + + EXPECT_NE(dynamic_cast(table.get()), nullptr) + << "Expected RestTable when server config requests server-side scan planning"; +} + +// -------------------------------------------------------------------------- +// Client config "scan-planning-mode":"server", table config absent → +// LoadTable returns a RestTable. +// -------------------------------------------------------------------------- +TEST_F(RestCatalogLoadTableTest, ClientConfigServerScanReturnsRestTable) { + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, kLoadTableDefaultResponse))); + + ICEBERG_UNWRAP_OR_FAIL(auto catalog, + MakeCatalog(all_plan_endpoints_, /*client_scan_mode=*/"server")); + ICEBERG_UNWRAP_OR_FAIL(auto as_catalog, catalog->AsCatalog()); + ICEBERG_UNWRAP_OR_FAIL(auto table, as_catalog->LoadTable(identifier_)); + + EXPECT_NE(dynamic_cast(table.get()), nullptr) + << "Expected RestTable when client config requests server-side scan planning"; +} + +// -------------------------------------------------------------------------- +// "scan-planning-mode":"server" but PlanTableScan endpoint missing → +// LoadTable returns NotSupported. +// -------------------------------------------------------------------------- +TEST_F(RestCatalogLoadTableTest, ServerScanWithMissingEndpointReturnsNotSupported) { + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, kLoadTableServerScanResponse))); + + ICEBERG_UNWRAP_OR_FAIL(auto catalog, MakeCatalog(no_plan_endpoint_set_)); + ICEBERG_UNWRAP_OR_FAIL(auto as_catalog, catalog->AsCatalog()); + auto result = as_catalog->LoadTable(identifier_); + + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); +} + +// -------------------------------------------------------------------------- +// No scan-planning-mode set anywhere → LoadTable returns a plain Table +// (not RestTable). +// -------------------------------------------------------------------------- +TEST_F(RestCatalogLoadTableTest, DefaultScanModeReturnsPlainTable) { + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, kLoadTableDefaultResponse))); + + ICEBERG_UNWRAP_OR_FAIL(auto catalog, MakeCatalog(all_plan_endpoints_)); + ICEBERG_UNWRAP_OR_FAIL(auto as_catalog, catalog->AsCatalog()); + ICEBERG_UNWRAP_OR_FAIL(auto table, as_catalog->LoadTable(identifier_)); + + EXPECT_EQ(dynamic_cast(table.get()), nullptr) + << "Expected plain Table (not RestTable) when no scan-planning-mode is configured"; +} + +} // namespace iceberg::rest diff --git a/src/iceberg/test/rest_json_serde_test.cc b/src/iceberg/test/rest_json_serde_test.cc index ec41e4a66..64e992d18 100644 --- a/src/iceberg/test/rest_json_serde_test.cc +++ b/src/iceberg/test/rest_json_serde_test.cc @@ -1569,13 +1569,6 @@ INSTANTIATE_TEST_SUITE_P( R"({"status":"submitted","plan-id":"somePlanId","plan-tasks":[]})", .expected_error_kind = ErrorKind::kValidationFailed, .expected_error_msg = "tasks can only be defined when status is 'completed'"}, - PlanTableScanResponseInvalidParam{ - .test_name = "FailedWithPlanId", - .json_str = - R"({"status":"failed","plan-id":"somePlanId","error":{"message":"x","type":"y","code":500}})", - .expected_error_kind = ErrorKind::kValidationFailed, - .expected_error_msg = - "plan id can only be defined when status is 'submitted' or 'completed'"}, PlanTableScanResponseInvalidParam{ .test_name = "FailedWithoutError", .json_str = R"({"status":"failed"})", @@ -2604,6 +2597,16 @@ TEST(PlanTableScanResponseRoundtripTest, FailedWithError) { EXPECT_EQ(*result, *result2); } +TEST(PlanTableScanResponseRoundtripTest, FailedWithPlanId) { + // A server may include a plan-id in a failed response so the client can cancel. + auto json = nlohmann::json::parse( + R"({"status":"failed","plan-id":"plan-123","error":{"message":"Planning failed","type":"PlanningException","code":500}})"); + auto result = PlanTableScanResponseFromJson(json, EmptySpecs(), EmptySchema()); + ASSERT_THAT(result, IsOk()); + EXPECT_EQ(result->plan_id, "plan-123"); + EXPECT_EQ(result->plan_status, PlanStatus::kFailed); +} + TEST(FetchPlanningResultResponseRoundtripTest, CompletedWithPlanTasks) { auto json = nlohmann::json::parse( R"({"status": "completed", "plan-tasks": ["task-1", "task-2"]})"); 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..e3c61e21c --- /dev/null +++ b/src/iceberg/test/rest_table_scan_test.cc @@ -0,0 +1,1089 @@ +/* + * 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/constants.h" +#include "iceberg/file_io.h" +#include "iceberg/manifest/manifest_entry.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; + +// Matches a JSON string body where `key` has exactly `expected_value`. +MATCHER_P2(JsonBodyHas, key, expected_value, "") { + try { + auto json = nlohmann::json::parse(arg); + if (!json.contains(key)) { + *result_listener << "JSON body missing key \"" << key << "\""; + return false; + } + nlohmann::json expected = expected_value; + if (json.at(key) != expected) { + *result_listener << "JSON[\"" << key << "\"] = " << json.at(key) << ", expected " + << expected; + return false; + } + return true; + } catch (...) { + *result_listener << "failed to parse JSON body"; + return false; + } +} + +// Matches a JSON string body that does NOT contain `key`. +MATCHER_P(JsonBodyLacks, key, "") { + try { + auto json = nlohmann::json::parse(arg); + if (json.contains(key)) { + *result_listener << "JSON body unexpectedly contains key \"" << key << "\""; + return false; + } + return true; + } catch (...) { + *result_listener << "failed to parse JSON body"; + return false; + } +} + +// -------------------------------------------------------------------------- +// 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)); +} + +// -------------------------------------------------------------------------- +// UseSnapshot(): the POST body sent to the server must contain both +// "snapshot-id" and "use-snapshot-schema": true. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, UseSnapshotPropagatesUseSnapshotSchemaInContext) { + constexpr int64_t kSnapshotId = 1000L; + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + + EXPECT_CALL(*mock_client_, + Post(_, + testing::AllOf(JsonBodyHas("snapshot-id", kSnapshotId), + JsonBodyHas("use-snapshot-schema", true)), + _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + RestTableScanBuilder builder(metadata_, file_io_, "test.my_table", nullptr, + MakeContext(std::nullopt)); + builder.UseSnapshot(kSnapshotId); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder.Build()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Default scan: POST body must have "use-snapshot-schema": false and no +// "snapshot-id" field. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, DefaultScanDoesNotSetUseSnapshotSchema) { + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + + EXPECT_CALL(*mock_client_, + Post(_, + testing::AllOf(JsonBodyHas("use-snapshot-schema", false), + JsonBodyLacks("snapshot-id")), + _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + RestTableScanBuilder builder(metadata_, file_io_, "test.my_table", nullptr, + MakeContext(std::nullopt)); + ICEBERG_UNWRAP_OR_FAIL(auto scan, builder.Build()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Storage credentials in COMPLETED response: 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()); + + // io() must return a credential-scoped IO without requiring a downcast. + EXPECT_NE(scan->io().get(), file_io_.get()); +} + +// -------------------------------------------------------------------------- +// No storage credentials: 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()); + + EXPECT_EQ(scan->io().get(), file_io_.get()); +} + +// -------------------------------------------------------------------------- +// Storage credentials returned in FetchScanTasksResponse also update 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()); + + EXPECT_NE(scan->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); +} + +// ========================================================================== +// PlanFilesStream tests +// ========================================================================== + +// -------------------------------------------------------------------------- +// PlanFilesStream: COMPLETED immediately, stream yields no tasks. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamCompleted) { + 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 stream, scan->PlanFilesStream()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, stream->ToVector()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: two plan-task tokens each trigger a separate FetchScanTasks +// POST, one per token, and the combined task set is returned. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamFetchesEachTokenSeparately) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-stream","plan-tasks":["tok-s1","tok-s2"]})"; + constexpr std::string_view kTask1Response = R"({ + "file-scan-tasks": [ + {"data-file":{"content":"data","file-path":"s3://b/f1.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}} + ] + })"; + constexpr std::string_view kTask2Response = R"({ + "file-scan-tasks": [ + {"data-file":{"content":"data","file-path":"s3://b/f2.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}} + ] + })"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kTask1Response)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kTask2Response)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto stream, scan->PlanFilesStream()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, stream->ToVector()); + ASSERT_EQ(tasks.size(), 2u); + EXPECT_EQ(tasks[0]->data_file()->file_path, "s3://b/f1.parquet"); + EXPECT_EQ(tasks[1]->data_file()->file_path, "s3://b/f2.parquet"); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: the second FetchScanTasks POST is not made until the first +// token's buffer is exhausted. Verified by counting POST calls between Next() +// invocations. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamFetchesTokenOnlyWhenBufferExhausted) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-lazy","plan-tasks":["tok-1","tok-2"]})"; + // tok-1 returns 2 tasks; tok-2 must not be fetched until both are consumed. + constexpr std::string_view kTwoTasksResponse = R"({ + "file-scan-tasks": [ + {"data-file":{"content":"data","file-path":"s3://b/f1.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}}, + {"data-file":{"content":"data","file-path":"s3://b/f2.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}} + ] + })"; + constexpr std::string_view kOneTaskResponse = R"({ + "file-scan-tasks": [ + {"data-file":{"content":"data","file-path":"s3://b/f3.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}} + ] + })"; + + int fetch_count = 0; + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce([&](auto&&...) -> Result { + ++fetch_count; + return HttpResponse::MakeForTesting(200, std::string(kTwoTasksResponse)); + }) + .WillOnce([&](auto&&...) -> Result { + ++fetch_count; + return HttpResponse::MakeForTesting(200, std::string(kOneTaskResponse)); + }); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto stream, scan->PlanFilesStream()); + + // No FetchScanTasks call yet — stream has not been driven. + EXPECT_EQ(fetch_count, 0); + + // First Next(): fetches tok-1 (2 tasks buffered), returns f1. + ICEBERG_UNWRAP_OR_FAIL(auto t1, stream->Next()); + ASSERT_TRUE(t1.has_value()); + EXPECT_EQ(fetch_count, 1); + EXPECT_EQ((*t1)->data_file()->file_path, "s3://b/f1.parquet"); + + // Second Next(): served from buffer; tok-2 not fetched yet. + ICEBERG_UNWRAP_OR_FAIL(auto t2, stream->Next()); + ASSERT_TRUE(t2.has_value()); + EXPECT_EQ(fetch_count, 1); + EXPECT_EQ((*t2)->data_file()->file_path, "s3://b/f2.parquet"); + + // Third Next(): buffer exhausted, fetches tok-2, returns f3. + ICEBERG_UNWRAP_OR_FAIL(auto t3, stream->Next()); + ASSERT_TRUE(t3.has_value()); + EXPECT_EQ(fetch_count, 2); + EXPECT_EQ((*t3)->data_file()->file_path, "s3://b/f3.parquet"); + + // Fourth Next(): all tokens consumed, stream terminates. + ICEBERG_UNWRAP_OR_FAIL(auto end, stream->Next()); + EXPECT_FALSE(end.has_value()); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: Next() propagates a FetchScanTasks error and DELETE /plan +// is called via the stream destructor since consumed_ is never set. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamNextReturnsErrorOnFetchFailure) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-next-err","plan-tasks":["tok-err"]})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(IOError("FetchScanTasks network error"))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto stream, scan->PlanFilesStream()); + auto result = stream->Next(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: stream destroyed after consuming the first task but before +// the second token is fetched → DELETE /plan called by the destructor. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamCancelAfterPartialConsumption) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-partial","plan-tasks":["tok-p1","tok-p2"]})"; + constexpr std::string_view kTaskResponse = R"({ + "file-scan-tasks": [ + {"data-file":{"content":"data","file-path":"s3://b/fp1.parquet", + "file-format":"PARQUET","spec-id":0,"partition":[],"file-size-in-bytes":1,"record-count":1}} + ] + })"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kPlanResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kTaskResponse)))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + { + ICEBERG_UNWRAP_OR_FAIL(auto stream, scan->PlanFilesStream()); + // Consume the first task from tok-p1; tok-p2 has never been fetched. + ICEBERG_UNWRAP_OR_FAIL(auto task, stream->Next()); + ASSERT_TRUE(task.has_value()); + // Destroy stream here — tok-p2 is still pending, so destructor calls DELETE. + } +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: stream destroyed with no Next() calls at all → +// DELETE /plan called by the destructor. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamCancelOnPartialConsumption) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"plan-never-consumed","plan-tasks":["tok-p1"]})"; + + 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())); + { + ICEBERG_UNWRAP_OR_FAIL(auto stream, scan->PlanFilesStream()); + // Destroy without any Next() call — destructor must call DELETE /plan. + } +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: SUBMITTED → poll until COMPLETED, stream yields all tasks. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamSubmittedThenCompleted) { + constexpr std::string_view kSubmittedBody = + R"({"status":"submitted","plan-id":"plan-poll-stream"})"; + 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 stream, scan->PlanFilesStream()); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, stream->ToVector()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: FAILED → stream returns an error. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamFailed) { + 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->PlanFilesStream(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// PlanFilesStream: FAILED with plan-id → DELETE /plan called before returning +// the error (tests the ExecuteScanPlanStream kFailed cancel fix). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, PlanFilesStreamFailedWithPlanIdCancels) { + constexpr std::string_view kFailedBody = + R"({"status":"failed","plan-id":"plan-fail-stream","error":{"message":"server error","type":"ServerError","code":500}})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kFailedBody)))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + auto result = scan->PlanFilesStream(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// FetchPlanningResult: GET fails after SUBMITTED → DELETE /plan called before +// returning the error (tests the FetchPlanningResult error-path cancel fix). +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, CancelCalledOnFetchPlanningResultGetError) { + constexpr std::string_view kSubmittedBody = + R"({"status":"submitted","plan-id":"plan-fetch-err"})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kSubmittedBody)))); + EXPECT_CALL(*mock_client_, Get(_, _, _, _, _)) + .WillOnce(Return(IOError("network failure"))); + 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)); +} + +// -------------------------------------------------------------------------- +// Stale credentials: second PlanFilesStream call resets scan_io_slot_ so that +// credentials from the first plan response do not persist into the second. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, SecondPlanFilesStreamCallClearsStaleCredentials) { + constexpr std::string_view kFirstResponse = R"({ + "status": "completed", + "storage-credentials": [ + {"prefix": "s3://bucket/prefix", "config": {"key": "value"}} + ] + })"; + constexpr std::string_view kSecondResponse = R"({"status":"completed"})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kFirstResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kSecondResponse)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeScan(MakeContext())); + + // First call: server vends credentials → io() returns a credential-scoped IO. + ICEBERG_UNWRAP_OR_FAIL(auto tasks1, scan->PlanFiles()); + EXPECT_TRUE(tasks1.empty()); + EXPECT_NE(scan->io().get(), file_io_.get()); + + // Second call: no credentials returned → io() must revert to the table IO, + // not retain the credentials from the first plan. + ICEBERG_UNWRAP_OR_FAIL(auto tasks2, scan->PlanFiles()); + EXPECT_TRUE(tasks2.empty()); + EXPECT_EQ(scan->io().get(), file_io_.get()); +} + +// ========================================================================== +// RestIncrementalAppendScan tests +// ========================================================================== + +class RestIncrementalAppendScanTest : public RestTableScanTest { + protected: + // Creates a RestIncrementalAppendScan with the given context and optional + // snapshot range. + Result> MakeIncrementalScan( + RestScanContext ctx, std::optional from_snapshot_id = std::nullopt, + bool from_inclusive = false, std::optional to_snapshot_id = std::nullopt) { + RestIncrementalAppendScanBuilder builder(metadata_, file_io_, "test.my_table", + nullptr, std::move(ctx)); + if (from_snapshot_id.has_value()) { + builder.FromSnapshot(*from_snapshot_id, from_inclusive); + } + if (to_snapshot_id.has_value()) { + builder.ToSnapshot(*to_snapshot_id); + } + return builder.Build(); + } +}; + +// -------------------------------------------------------------------------- +// PlanFiles: server returns COMPLETED immediately, no tasks. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, 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, MakeIncrementalScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: COMPLETED with a plan-task token; FetchScanTasks is called. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesWithPlanTasks) { + constexpr std::string_view kPlanResponse = + R"({"status":"completed","plan-id":"incr-plan-1","plan-tasks":["tok-incr-1"]})"; + 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, MakeIncrementalScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: SUBMITTED → poll → COMPLETED. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesSubmittedThenCompleted) { + constexpr std::string_view kSubmittedBody = + R"({"status":"submitted","plan-id":"incr-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, MakeIncrementalScan(MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: FAILED → IOError. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, 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, MakeIncrementalScan(MakeContext())); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// PlanFiles: PlanTableScan endpoint missing → NotSupported. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesEndpointNotSupported) { + ICEBERG_UNWRAP_OR_FAIL( + auto scan, MakeIncrementalScan(MakeContext(std::unordered_set{}))); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kNotSupported)); +} + +// -------------------------------------------------------------------------- +// No current snapshot → PlanFiles returns empty without calling the server. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesEmptyWhenNoCurrentSnapshot) { + // Build metadata with no current snapshot. + auto spec = PartitionSpec::Unpartitioned(); + auto empty_metadata = std::make_shared(TableMetadata{ + .format_version = 2, + .table_uuid = "no-snap-uuid", + .location = "/tmp/table", + .last_sequence_number = 0L, + .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 = kInvalidSnapshotId, + }); + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)).Times(0); + + // Use Make() directly to bypass builder validation that requires a snapshot. + ICEBERG_UNWRAP_OR_FAIL(auto scan, RestIncrementalAppendScan::Make( + empty_metadata, schema_, file_io_, + internal::TableScanContext{}, MakeContext())); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Explicit to_snapshot_id: POST body must contain "end-snapshot-id" set to +// the given value, not the current table snapshot. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesWithExplicitToSnapshotId) { + constexpr int64_t kToSnapshotId = 1000L; + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, + Post(_, JsonBodyHas("end-snapshot-id", kToSnapshotId), _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL( + auto scan, MakeIncrementalScan(MakeContext(), std::nullopt, false, kToSnapshotId)); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Exclusive from_snapshot_id: POST body must pass from_snapshot_id directly +// as "start-snapshot-id" (exclusive), and current snapshot as "end-snapshot-id". +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesWithFromSnapshotIdExclusive) { + constexpr int64_t kFromSnapshotId = 999L; + constexpr int64_t kCurrentSnapshotId = 1000L; + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, + Post(_, + testing::AllOf(JsonBodyHas("start-snapshot-id", kFromSnapshotId), + JsonBodyHas("end-snapshot-id", kCurrentSnapshotId)), + _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeIncrementalScan(MakeContext(), kFromSnapshotId, + /*inclusive=*/false)); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// Inclusive from_snapshot_id: parent snapshot is used as start_snapshot_id. +// The fixture snapshot (id=1000) has no parent, so "start-snapshot-id" is +// absent from the POST body. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesWithFromSnapshotIdInclusiveNoParent) { + constexpr int64_t kFromSnapshotId = 1000L; + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, + Post(_, + testing::AllOf(JsonBodyLacks("start-snapshot-id"), + JsonBodyHas("end-snapshot-id", kFromSnapshotId)), + _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL( + auto scan, MakeIncrementalScan(MakeContext(), kFromSnapshotId, /*inclusive=*/true)); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// -------------------------------------------------------------------------- +// PlanFiles: FAILED with plan-id → DELETE /plan called before returning the +// error (tests the ExecuteScanPlan kFailed cancel fix). +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesFailedWithPlanIdCancels) { + constexpr std::string_view kFailedBody = + R"({"status":"failed","plan-id":"plan-fail-incr","error":{"message":"server error","type":"ServerError","code":500}})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kFailedBody)))); + EXPECT_CALL(*mock_client_, Delete(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, "{}"))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeIncrementalScan(MakeContext())); + auto result = scan->PlanFiles(); + EXPECT_THAT(result, IsError(ErrorKind::kIOError)); +} + +// -------------------------------------------------------------------------- +// Stale credentials: second PlanFiles call resets scan_io_ so that credentials +// from the first plan response do not persist into the second. +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, SecondPlanFilesCallClearsStaleCredentials) { + constexpr std::string_view kFirstResponse = R"({ + "status": "completed", + "storage-credentials": [ + {"prefix": "s3://bucket/prefix", "config": {"key": "value"}} + ] + })"; + constexpr std::string_view kSecondResponse = R"({"status":"completed"})"; + + EXPECT_CALL(*mock_client_, Post(_, _, _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kFirstResponse)))) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kSecondResponse)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeIncrementalScan(MakeContext())); + + // First call: server vends credentials. + ICEBERG_UNWRAP_OR_FAIL(auto tasks1, scan->PlanFiles()); + EXPECT_TRUE(tasks1.empty()); + + // Second call: no credentials. Verifies scan_io_ was cleared so stale + // credentials from the first plan do not bleed into the second request. + ICEBERG_UNWRAP_OR_FAIL(auto tasks2, scan->PlanFiles()); + EXPECT_TRUE(tasks2.empty()); +} + +// -------------------------------------------------------------------------- +// Inclusive from_snapshot_id with a parent: POST body must use the parent's id +// as "start-snapshot-id" and the current snapshot as "end-snapshot-id". +// -------------------------------------------------------------------------- +TEST_F(RestIncrementalAppendScanTest, PlanFilesWithFromSnapshotIdInclusiveWithParent) { + constexpr int64_t kParentSnapshotId = 900L; + constexpr int64_t kChildSnapshotId = 1001L; + constexpr int64_t kCurrentSnapshotId = 1000L; + + // Add a second snapshot with a parent to the metadata. + auto child_snapshot = std::make_shared( + Snapshot{.snapshot_id = kChildSnapshotId, + .parent_snapshot_id = kParentSnapshotId, + .sequence_number = 2L, + .timestamp_ms = TimePointMsFromUnixMs(1609459260000L), + .manifest_list = "/tmp/manifest-list-2.avro", + .schema_id = schema_->schema_id()}); + metadata_->snapshots.push_back(child_snapshot); + + constexpr std::string_view kResponseBody = R"({"status":"completed"})"; + EXPECT_CALL(*mock_client_, + Post(_, + testing::AllOf(JsonBodyHas("start-snapshot-id", kParentSnapshotId), + JsonBodyHas("end-snapshot-id", kCurrentSnapshotId)), + _, _, _)) + .WillOnce(Return(HttpResponse::MakeForTesting(200, std::string(kResponseBody)))); + + ICEBERG_UNWRAP_OR_FAIL(auto scan, MakeIncrementalScan(MakeContext(), kChildSnapshotId, + /*inclusive=*/true)); + ICEBERG_UNWRAP_OR_FAIL(auto tasks, scan->PlanFiles()); + EXPECT_TRUE(tasks.empty()); +} + +// ========================================================================== +// RestTable::NewIncrementalAppendScan +// ========================================================================== + +// -------------------------------------------------------------------------- +// RestTable::NewIncrementalAppendScan returns a RestIncrementalAppendScanBuilder. +// -------------------------------------------------------------------------- +TEST_F(RestTableScanTest, RestTableNewIncrementalAppendScanReturnsRestBuilder) { + 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->NewIncrementalAppendScan()); + 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 c9af765e9..8deada17d 100644 --- a/src/iceberg/test/table_scan_test.cc +++ b/src/iceberg/test/table_scan_test.cc @@ -841,6 +841,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