Skip to content

Ingestion client: ingest, stream, request status (#17, PR A) - #67

Merged
craigmcchesney merged 4 commits into
mainfrom
feat/17-ingestion-client
Sep 30, 2026
Merged

craigmcchesney merged 4 commits into
mainfrom
feat/17-ingestion-client

Conversation

@craigmcchesney

Copy link
Copy Markdown
Collaborator

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)

  • Builders for the column kinds interface to modernized annotation API #6 left out: double_array_column() and its float/int32/int64/bool siblings (nested sequences or NumPy arrays, flattened row-major, 1–3 dims, optional dims=), 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 under max_rows / max_bytes / max_span_nanos. max_bytes bounds the whole IngestDataRequest, including two worst-case ids; a test checks the arithmetic against protobuf's own ByteSize(). 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, required since).
  • _dispatch gains an optional rpc_error_hint; a size-violation RESOURCE_EXHAUSTED gets a pointer to split_data_frame() (or, on a status query, to a narrower query). Other RESOURCE_EXHAUSTED errors 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.md section.

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.py task.

  1. Requests sent before a producer failure are not necessarily delivered. The plan (T10/D5) said requests before the failing one "have already been ingested". The in-process grpcio tests showed grpcio's cancellation racing the sends: one of two requests handed over never reached the server. The re-raised exception's note now says "any that reached the server are ingested — check request status".
  2. The sent count travels as an exception note, not in the message. Putting it in the message would mean changing the exception's type or mutating its args. So it uses a PEP 678 note, Python 3.11+ only, and is always logged.
  3. Abandoning a bidi stream now cancels the call. The plan didn't anticipate this. Without it, grpcio keeps pulling and sending the caller's requests on its own thread after the caller stops reading. An in-process test saw all 1,000 queued requests drain into the server, and that test fails with the fix removed.
  4. data_frame() doesn't call validate_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 raises ValueError for hand-built columns the server rejects anyway: an enum without enumId, an array with >3 dims, an image without a complete descriptor, a struct without schemaId, a serialized column without encoding. enum_column() rejects a whitespace-only enum_id. The three existing tests that hand-built such columns (T12) were fixed, not deleted. All of this is in NEXT.md.

Testing

  • 873 unit tests pass locally on Python 3.12. The 3.10–3.13 matrix runs in CI. The only version-dependent path is the exception note, which is hasattr-guarded, with a version-gated assertion.
  • tests/unit/test_ingestion_streaming_grpc.py runs the streaming paths against a real in-process grpcio server on localhost:0. It passed 25 consecutive runs. It also pins the grpcio premise itself: through the raw stub, the caller's exception is lost.
  • ruff lint and format, mypy src/, and the cookbook and release-notes checkers are clean.
  • mypy against typed stubs: I generated .pyi stubs locally with dp-grpc's pinned grpcio-tools==1.84.0 / mypy-protobuf==5.1.0 and its generate-python-stubs.yml flags, per CLAUDE.md. This caught one enum-typing error, now fixed. The CLAUDE.md snippet was type-checked against src/ separately, since the cookbook checker does not cover CLAUDE.md.
  • No live-server testing in this PR. That is PR B.

🤖 Generated with Claude Code

https://claude.ai/code/session_01QkZu3MPfUopTisahmvhakD

craigmcchesney and others added 3 commits September 30, 2026 13:58
- 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
Copilot AI balanced review requested due to automatic review settings September 30, 2026 20:19

Copilot AI left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 High severity · 4 Low severity

Open (5)
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.

Comment thread src/dp_python_lib/client/ingestion_client.py
Comment thread CLAUDE.md Outdated
Comment thread doc/release-notes/NEXT.md Outdated
Comment thread plan/tickets/17/plan.md Outdated
Comment thread src/dp_python_lib/client/ingestion_client.py
)

- 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
@craigmcchesney
craigmcchesney merged commit 92c3882 into main Sep 30, 2026
6 checks passed
@craigmcchesney
craigmcchesney deleted the feat/17-ingestion-client branch September 30, 2026 21:05
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants