Repository navigation
feat(rest): add scan plan endpoint support to REST catalog client - #783
gsandeep1241 wants to merge 17 commits into
Conversation
|
This pull request has been marked as stale due to 30 days of inactivity. It will be closed in 1 week if no further activity occurs. If you think that’s incorrect or this pull request requires a review, please simply write any comment. If closed, you can revive the PR at any time and @mention a reviewer or discuss it on the dev@iceberg.apache.org list. Thank you for your contributions. |
|
This pull request has been closed due to lack of activity. This is not a judgement on the merit of the PR in any way. It is just a way of keeping the PR queue manageable. If you think that is incorrect, or the pull request requires review, you can revive the PR at any time. |
|
Hey @gsandeep1241, do you want to revive this? |
|
@wgtmac Thanks for re-opening it. Apologies for the long delay, I'll get back on it this week - next set of changes should be out for review in the next couple of days! |
fb8da85 to
d8a7d6c
Compare
2bab861 to
bcc9da4
Compare
Thanks @wgtmac for reviving this! This is now ready for review. Please take a look when you can :) |
When a table is loaded from a REST catalog that advertises the PlanTableScan
endpoint, NewScan() now returns a RestTableScanBuilder whose Build() produces
a RestTableScan. PlanFiles() on that scan delegates manifest resolution to
the server via POST /plan, GET /plan/{id} (with exponential backoff),
POST /tasks/{id}, and DELETE /plan/{id} (best-effort cancel), instead of
reading manifests locally.
- Add RestTable, RestTableScanBuilder, RestTableScan and RestScanContext
- Promote DataTableScan::PlanFiles and TableScanBuilder::Build to virtual
- Convert RestCatalog::client_ and paths_ to shared_ptr so RestScanContext
can share ownership with live scans
- Cancel server-side plan when ResolveScanTasks fails partway through - Propagate use_snapshot_schema from scan context to PlanTableScanRequest: true for UseSnapshot/AsOfTime/tag refs and incremental scans, false for branch refs and default scans - Gate RestTable creation on effective scan-planning-mode config (table config overrides client config, default is client); error if server mode is requested but endpoint is not advertised - Add ScanPlanningMode enum and ScanPlanningModeFrom() parser to RestCatalogProperties - Make HttpClient methods virtual and add HttpResponse::MakeForTesting() to support unit test mocking - Add tests: use_snapshot_schema in table_scan_test, ScanPlanningModeFrom parsing in catalog_properties_test, and RestTableScan HTTP flow tests in rest_table_scan_test
Upstream changed Table and DataTableScanBuilder constructors to require full_name/table_name and MetricsReporter parameters. Updated RestTable, RestTableScanBuilder, and their callers accordingly.
Parse storage-credentials from PlanTableScanResponse, FetchPlanningResultResponse, and FetchScanTasksResponse. When credentials are present, build a scan-scoped FileIO via MakeTableFileIO and expose it through RestTableScan::effective_io() for callers to use when reading the returned scan tasks.
RestTableScanBuilder::Build() calls context_.Validate() across the iceberg_rest/iceberg library boundary. Without ICEBERG_EXPORT on TableScanContext the symbol is hidden in the shared library and the linker fails on arm64.
…port MSVC does not export the implicitly-generated move constructor of a template class instantiation. RestTableScanBuilder (introduced in iceberg_rest) is exported with ICEBERG_REST_EXPORT, so its compiler- generated move constructor must call the base TableScanBuilder move constructor as an imported symbol. Explicitly defaulting it makes it part of the explicit template instantiation and therefore exported.
…lesStream, and virtual io() - Add RestIncrementalAppendScan and RestIncrementalAppendScanBuilder; override RestTable::NewIncrementalAppendScan() to delegate PlanFiles() to the REST server via the same scan planning endpoints. IncrementalChangelogScan is left local: the planTableScan response only carries FileScanTask objects, not the per-snapshot operation metadata needed to reconstruct ChangelogScanTask entries. Refactor shared HTTP planning logic into free functions (ExecuteScanPlan, FetchPlanningResult, FetchScanTasks, ResolveScanTasks, ApplyStorageCredentials, CancelPlanning) used by both scan types. - Make DataTableScan::PlanFilesStream() virtual and implement it in RestTableScan with a lazy RestFileScanTaskStream that drives FetchScanTasks per plan_task token. The base-class PlanFiles() delegates to PlanFilesStream()->ToVector(), so the server planning path is used for both the eager and streaming APIs. - Make TableScan::io() virtual and override it in RestTableScan to return the credential-scoped FileIO when the server vends storage credentials, without requiring a downcast. Remove effective_io(). A shared ScanIoSlot (shared_ptr<shared_ptr<FileIO>>) is held by both the scan and the lazy stream so credentials vended by any FetchScanTasks response are visible through io() even after the stream is consumed. Update tests accordingly.
a9cbb7d to
5de0344
Compare
PlanFiles() was made virtual in this PR to support a RestTableScan override. That override has since been replaced by RestTableScan::PlanFilesStream(), so PlanFiles() no longer needs to be virtual — callers dispatch through the virtual PlanFilesStream() hook instead.
… and virtual io() - PlanFilesStream: COMPLETED, lazy token fetching, cancel-on-partial-consume, SUBMITTED→poll→COMPLETED, FAILED - RestIncrementalAppendScan: COMPLETED, plan-tasks, SUBMITTED→poll→COMPLETED, FAILED, endpoint-not-supported, no-current-snapshot returns empty, explicit to_snapshot_id, exclusive from_snapshot_id, inclusive from_snapshot_id with and without a parent snapshot - RestTable::NewIncrementalAppendScan returns RestIncrementalAppendScanBuilder
FetchPlanningResult returned early on GET failure, JSON parse error, response deserialization error, Validate() failure, and ApplyStorageCredentials failure without cancelling, leaving the server holding plan resources. ExecuteScanPlanStream and ExecuteScanPlan had the same gap for the kFailed status and ApplyStorageCredentials failure in the kCompleted branch. Replace ICEBERG_ASSIGN_OR_RAISE / ICEBERG_RETURN_UNEXPECTED with explicit checks in FetchPlanningResult's poll loop, and add CancelPlanning before every error return in ExecuteScanPlanStream and ExecuteScanPlan where plan_id is known.
If PlanFilesStream()/PlanFiles() is called more than once and only the first response vends storage credentials, the second plan would inherit the stale credential-scoped FileIO from the first. Reset the slot to null before each planning call so that io() falls back to the table IO when the server does not vend credentials in the new response.
…ial fixes - Add JsonBodyHas/JsonBodyLacks GMock matchers that parse the POST body so request field propagation is verified against actual server JSON, not just internal context state - Update UseSnapshot and DefaultScan tests to call PlanFiles() and match "snapshot-id"/"use-snapshot-schema" in the POST body - Update incremental scan snapshot routing tests to match "start-snapshot-id"/"end-snapshot-id" values in the POST body - Add PlanFilesStreamFailedWithPlanIdCancels: verifies DELETE is called when ExecuteScanPlanStream receives FAILED with a plan-id - Add CancelCalledOnFetchPlanningResultGetError: verifies DELETE is called when FetchPlanningResult GET fails after SUBMITTED - Add SecondPlanFilesStreamCallClearsStaleCredentials: verifies scan_io_slot_ is reset so second plan response with no credentials reverts io() to table IO - Add PlanFilesFailedWithPlanIdCancels for RestIncrementalAppendScan - Add SecondPlanFilesCallClearsStaleCredentials for RestIncrementalAppendScan
…partial-consume tests - Rename PlanFilesStreamLazilyFetchesPlanTasks to PlanFilesStreamFetchesEachTokenSeparately; verify task content and count rather than empty responses, since the old comment implied laziness was being tested when only call-count was verified - Add PlanFilesStreamFetchesTokenOnlyWhenBufferExhausted: uses a counter incremented inside WillOnce lambdas to assert the second FetchScanTasks POST is not made until the first token buffer is fully drained - Add PlanFilesStreamNextReturnsErrorOnFetchFailure: calls Next() directly and asserts it returns IOError when FetchScanTasks fails - Add PlanFilesStreamCancelAfterPartialConsumption: consumes one task then destroys the stream, verifying DELETE called with a token still pending - Retain PlanFilesStreamCancelOnPartialConsumption with updated comment clarifying it tests destroy-before-any-Next
…ests
- Add MakeForTesting() factory to RestCatalog that bypasses FetchServerConfig,
enabling unit tests to inject pre-built client/paths/endpoints.
- Add rest_catalog_unit_test.cc with four tests covering:
table config server scan -> RestTable, client config server scan -> RestTable,
server scan with missing PlanTableScan endpoint -> NotSupported,
and default (no config) -> plain Table.
- Fix RestIncrementalAppendScan::PlanFiles() to check kInvalidSnapshotId
before calling Snapshot(), so a table with no current snapshot returns
empty results rather than a NotFound error.
- Fix make_unique<RestFileScanTaskStream> deduction failure by replacing
brace-init {} with std::vector<std::string>{} for plan_task_tokens arg.
- Relax PlanTableScanResponse::Validate() to allow plan-id in failed
responses, enabling the cancel-on-fail code path to be reached.
- Update rest_json_serde_test.cc: remove FailedWithPlanId invalid case,
add FailedWithPlanId roundtrip test.
- Add iceberg/manifest/manifest_entry.h and iceberg/constants.h includes
to rest_table_scan_test.cc to fix DataFile and kInvalidSnapshotId access.
|
@wgtmac Thank you for the detailed review. I've addressed your comments. Can you please take a look once more when you have the chance? |
| 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; |
There was a problem hiding this comment.
When no explicit to_snapshot_id is set, this always uses the table's current snapshot and ignores context_.branch. A branch whose head differs from main will send the wrong end-snapshot-id; please resolve the branch head here and validate explicit endpoints against it.
| RestScanContext rest_context); | ||
|
|
||
| RestScanContext rest_context_; | ||
| mutable std::shared_ptr<FileIO> scan_io_; |
There was a problem hiding this comment.
This stores the vended scan IO, but the incremental scan does not override TableScan::io(). Normal callers therefore keep using the table IO after PlanFiles(), so credentials returned for incremental planning are never used to read the tasks. Please expose the effective IO here as well.
| Result<std::unique_ptr<DataTableScan>> RestTableScanBuilder::Build() { | ||
| ICEBERG_RETURN_UNEXPECTED(CheckErrors()); | ||
| ICEBERG_RETURN_UNEXPECTED(context_.Validate()); | ||
| ICEBERG_ASSIGN_OR_RAISE(auto schema, ResolveSnapshotSchema()); |
There was a problem hiding this comment.
ResolveSnapshotSchema() always returns the snapshot schema when snapshot_id is set, even when UseRef() marked a branch scan with use_snapshot_schema=false. Java uses the current table schema for branches and the snapshot schema for tags/time travel; the schema passed to this REST scan should follow that flag.
| .catalog_config = config_.configs(), | ||
| .table_config = table_config, | ||
| }; | ||
| return RestTable::Make(identifier, std::move(result.metadata), |
There was a problem hiding this comment.
This creates a RestTable only when the table is loaded. MakeTableFromCommitResponse() still returns a plain Table, so committing through a server-planning table silently switches later scans back to local manifest planning. Please preserve the REST table type and scan context after commit.
| std::shared_ptr<FileIO>& scan_io) { | ||
| if (credentials.empty()) return {}; | ||
| ICEBERG_ASSIGN_OR_RAISE( | ||
| auto io, MakeTableFileIO(ctx.catalog_config, ctx.table_config, credentials)); |
There was a problem hiding this comment.
This builds the scan IO from static credentials only. It does not install a StorageCredentialProvider or pass the scan plan-id; Java adds rest-scan-plan-id and refreshes vended credentials through the credentials endpoint. Temporary credentials can otherwise expire during a scan.
| ICEBERG_UNWRAP_OR_FAIL(auto tasks1, scan->PlanFiles()); | ||
| EXPECT_TRUE(tasks1.empty()); | ||
|
|
||
| // Second call: no credentials. Verifies scan_io_ was cleared so stale |
There was a problem hiding this comment.
These assertions only check that planning returns tasks. Please also call scan->io() after vended credentials are returned; the public incremental-scan path currently exposes the original table IO.
| // UseSnapshot(): the POST body sent to the server must contain both | ||
| // "snapshot-id" and "use-snapshot-schema": true. | ||
| // -------------------------------------------------------------------------- | ||
| TEST_F(RestTableScanTest, UseSnapshotPropagatesUseSnapshotSchemaInContext) { |
There was a problem hiding this comment.
Please add request-body tests for Project() and IncludeColumnStats({"..."}). The current suite checks snapshot flags, but it does not catch the missing select projection or stats-fields described above.
| // -------------------------------------------------------------------------- | ||
| // PlanFilesStream: SUBMITTED → poll until COMPLETED, stream yields all tasks. | ||
| // -------------------------------------------------------------------------- | ||
| TEST_F(RestTableScanTest, PlanFilesStreamSubmittedThenCompleted) { |
There was a problem hiding this comment.
This submitted case has no plan tasks, so it cannot verify lazy fetching. Please return multiple plan-task tokens from the poll response and assert that /tasks is not called until Next() consumes the stream.
| EXPECT_FALSE(scan->context().use_snapshot_schema); | ||
| } | ||
|
|
||
| INSTANTIATE_TEST_SUITE_P(TableScanVersions, TableScanTest, testing::Values(1, 2, 3)); |
There was a problem hiding this comment.
These use_snapshot_schema checks do not depend on format version, but the parameterization runs each one for versions 1, 2, and 3. Please make these non-parameterized or avoid repeating the same assertions three times.
| // PlanFilesStream: stream destroyed with no Next() calls at all → | ||
| // DELETE /plan called by the destructor. | ||
| // -------------------------------------------------------------------------- | ||
| TEST_F(RestTableScanTest, PlanFilesStreamCancelOnPartialConsumption) { |
There was a problem hiding this comment.
The test name says partial consumption, but the stream is destroyed without any Next() call. Please rename it to make the scenario clear, such as CancelOnNoConsumption.
When a table is loaded from a REST catalog that advertises the PlanTableScan endpoint, NewScan() now returns a RestTableScanBuilder whose Build() produces a RestTableScan. PlanFiles() on that scan delegates manifest resolution to the server via POST /plan, GET /plan/{id} (with exponential backoff), POST /tasks/{id}, and DELETE /plan/{id} (best-effort cancel), instead of reading manifests locally.