Ingestion client: ingest, stream, request status (#17, PR A) - #67
Merged
Merged
Conversation
- Non-scalar column builders: double/float/int32/int64/bool_array_column (shape inferred from nested sequences or NumPy arrays, row-major, 1-3 dims, optional explicit dims), image_column, struct_column, serialized_column. - _check_column() enforces each kind's structural fields on hand-built columns too (enumId, 1-3 array dims, image descriptor, schemaId, serialized encoding): T7. validate_data_frame() applies data_frame()'s checks to an assembled frame. - split_data_frame(): lazy time-axis chunking under max_rows / max_bytes / max_span_nanos. max_bytes bounds the whole IngestDataRequest with room for two MAX_BUDGETED_ID_CHARS ids; SamplingClock chunk starts are integer nanoseconds. SERVER_DEFAULT_MAX_MESSAGE_BYTES is exported as a reference. - time_conversions.from_epoch_nanos() replaces two private copies, ahead of split_data_frame() needing a third. - Fix the three test_data_frame.py tests that hand-built columns the new checks reject (T12). Refs #17 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD
- IngestDataRequestParams (generated uuid4 request ids, ids capped at MAX_BUDGETED_ID_CHARS, frame re-validated) and chunked_request_params() for split_data_frame() chunks named <base>-<n>. - ingest_data() via _dispatch. ingest_data_stream() is hand-written so a partial reject keeps its response and rejected_request_ids. iter_ingest_data_bidi_stream() yields per-request rejects rather than raising, raises RuntimeError on transport errors, and cancels the call when the caller stops reading. - _RequestFeed re-raises the caller's own request-iterator exception in place of grpcio's "Exception iterating requests!", chained from the RpcError, with a sent-count note (3.11+) and log line. - RequestStatusQuery (RS), IngestionRequestStatus, query_request_status(), and await_request_statuses(): one provider + time-range query per poll, ids matched client-side, a required `since` floor backed off by REQUEST_STATUS_CLOCK_SKEW. - _dispatch gains an rpc_error_hint hook; size-violation RESOURCE_EXHAUSTED errors get a split_data_frame() or narrow-the-query hint, other RESOURCE_EXHAUSTED errors keep the plain message. - RegisterProviderApiResult.provider_id / .is_new_provider. - In-process grpcio tests for the streaming paths; plan corrected where they showed delivery before a failure is not guaranteed. Refs #17 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD
- CLAUDE.md: Key Files entries for ingestion_client.py, the two new test modules, and the extended data_frame.py / time_conversions.py; the rpc_error_hint hook and the two hand-written ingest senders in the client pattern; an Ingestion API section with the invariants that outlive the ticket (ack is not ingestion, SUCCESS is the zero value, one-id criteria and no paging behind the await design, the stream-keeps-response rule, the re-raise and its delivery caveat, cancel on abandon, the size caps). - README.md: ingestion moves from TODO to implemented, except subscribeData. - NEXT.md: a section for #17, including the new data_frame() rejections of hand-built columns as a behavior change. Refs #17 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Early bidi-loop exit does not reliably cancel a retained generator as documented, and related error branches need coverage.
Review effort: Balanced
Findings: 1
Open (5)
Do not rely on break for deterministic generator cancellation · New Narrow cancellation guarantee to explicit or context-managed close · New Correct scope of wrapped ingestion API operations · New Document explicit-close behavior for retained generators · New Add coverage for unexpected stream sender exceptions · New
What changed in this PR
Adds ingestion, streaming, request-status polling, DataFrame chunking, and non-scalar column support.
Changes:
- Implements unary and streaming ingestion APIs with status polling.
- Adds validated non-scalar columns and size/span-based frame chunking.
- Expands exports, documentation, release notes, and tests.
| File | Description |
|---|---|
src/dp_python_lib/client/ingestion_client.py |
Implements ingestion and status APIs. |
src/dp_python_lib/client/data_frame.py |
Adds builders, validation, and chunking. |
src/dp_python_lib/client/time_conversions.py |
Adds nanoseconds-to-timestamp conversion. |
src/dp_python_lib/client/service_api_client_base.py |
Supports contextual gRPC error hints. |
src/dp_python_lib/client/sample_status_conversions.py |
Reuses the shared timestamp converter. |
src/dp_python_lib/client/data_frame_conversions.py |
Reuses the shared timestamp converter. |
src/dp_python_lib/client/__init__.py |
Exports the new public APIs. |
tests/unit/test_ingestion_client.py |
Tests ingestion and polling behavior. |
tests/unit/test_ingestion_streaming_grpc.py |
Tests real grpcio streaming behavior. |
tests/unit/test_data_frame.py |
Tests builders, validation, and chunking. |
tests/unit/test_data_frame_conversions.py |
Tests non-scalar read-back symmetry. |
tests/unit/test_time_conversions.py |
Tests inverse timestamp conversion. |
tests/unit/test_service_api_client_base.py |
Tests gRPC error hints. |
README.md |
Documents ingestion support. |
CLAUDE.md |
Records ingestion architecture and usage. |
plan/tickets/17/plan.md |
Records implementation departures. |
doc/release-notes/NEXT.md |
Adds user-visible release notes. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
) - Bidi cancellation: closing the iterator cancels the call, but a bare break on a still-referenced generator does not close it (an in-process run saw all 1,000 queued requests drain). The docstring, CLAUDE.md, NEXT.md, and the plan now say to stop early with contextlib.closing; tests pin both the documented pattern and the bare-break premise. - validate_data_frame() takes a caller= name for its messages, so frames rejected by IngestDataRequestParams or split_data_frame() no longer report as data_frame() errors. - chunked_request_params() checks base_request_id (non-blank, room for the "-<n>" suffix) before consuming any frame. - The oversized-request hint takes its byte count from SERVER_DEFAULT_MAX_MESSAGE_BYTES instead of repeating the number. - IngestDataRequestParams' docstring notes that lazily built params fail mid-stream. - NEXT.md no longer says the whole ingestion API is covered; subscribeData() is not wrapped. - Unit coverage for the unexpected-exception tier of both hand-written streaming senders. Refs #17 Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.


Refs #17
PR A of the two planned in
plan/tickets/17/plan.md: the ingestion client and its unit tests. PR B (live tests, the cookbook recipe, and moving the integration tests' stub-based ingest onto the client) follows and closes the issue.What's in it
data_frame.py(69a74af)double_array_column()and its float/int32/int64/bool siblings (nested sequences or NumPy arrays, flattened row-major, 1–3 dims, optionaldims=),image_column(),struct_column(),serialized_column()._check_column()now enforces each kind's structural fields on hand-built columns too (enumId, 1–3 dims, image descriptor, schemaId, serialized encoding): T7.validate_data_frame()applies the same checks to an assembled frame.split_data_frame(): lazy time-axis chunking undermax_rows/max_bytes/max_span_nanos.max_bytesbounds the wholeIngestDataRequest, including two worst-case ids; a test checks the arithmetic against protobuf's ownByteSize(). A 2M-row frame splits in 0.15 s, with the largest chunk at 4,095,995 of 4,096,000 bytes.time_conversions.from_epoch_nanos()replaces two private copies, ahead of a third.ingestion_client.py(1bb1a20)IngestDataRequestParams,chunked_request_params(),ingest_data(),ingest_data_stream(),iter_ingest_data_bidi_stream().RequestStatusQuery(RS),IngestionRequestStatus,query_request_status(),await_request_statuses()(one provider + time-range query per poll, ids matched client-side, requiredsince)._dispatchgains an optionalrpc_error_hint; a size-violationRESOURCE_EXHAUSTEDgets a pointer tosplit_data_frame()(or, on a status query, to a narrower query). OtherRESOURCE_EXHAUSTEDerrors keep the plain message.RegisterProviderApiResult.provider_id/.is_new_provider.Docs (094df3c): CLAUDE.md (Key Files, the client pattern, an Ingestion API section), README, and a
NEXT.mdsection.Where the implementation departs from the plan
The plan is updated in this PR to match; see the dated notes under D5 and the
data_frame.pytask.data_frame()doesn't callvalidate_data_frame()after assembly. Both share one column-list check.data_frame()runs it on the caller's list, so an unsupported type is caught before it has to be routed, and messages keep the caller's index.Behavior change
data_frame()now raisesValueErrorfor hand-built columns the server rejects anyway: an enum withoutenumId, an array with >3 dims, an image without a complete descriptor, a struct withoutschemaId, a serialized column withoutencoding.enum_column()rejects a whitespace-onlyenum_id. The three existing tests that hand-built such columns (T12) were fixed, not deleted. All of this is inNEXT.md.Testing
hasattr-guarded, with a version-gated assertion.tests/unit/test_ingestion_streaming_grpc.pyruns the streaming paths against a real in-process grpcio server onlocalhost:0. It passed 25 consecutive runs. It also pins the grpcio premise itself: through the raw stub, the caller's exception is lost.mypy src/, and the cookbook and release-notes checkers are clean..pyistubs locally with dp-grpc's pinnedgrpcio-tools==1.84.0/mypy-protobuf==5.1.0and itsgenerate-python-stubs.ymlflags, per CLAUDE.md. This caught one enum-typing error, now fixed. The CLAUDE.md snippet was type-checked againstsrc/separately, since the cookbook checker does not cover CLAUDE.md.🤖 Generated with Claude Code
https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD