You signed in with another tab or window. Reload to refresh your session.You signed out in another tab or window. Reload to refresh your session.You switched accounts on another tab or window. Reload to refresh your session.Dismiss alert
Triaged 2026-09-27. This was an AI-drafted ticket. The scope below replaces the original draft,
which is kept at the bottom. The full plan, with server behavior verified against dp-service main,
is in plan/tickets/17/plan.md (plan PR #66).
Summary
Add the ingestion side of the Python client. Today IngestionClient wraps only registerProvider().
Without an ingestion client, the closed-loop query test (#7), the cookbook's ingestion recipe, and the
live verification of sample-status filtering are all blocked.
What triage changed
The "shared payload/params model" already exists.interface to modernized annotation API #6 built data_frame.py / data_frame_conversions.py
as the substrate for this ticket (ingestionDataFrame is the same common.DataFrame as calculations).
So Phase 1 is a thin IngestDataRequestParams(provider_id, frame, client_request_id) plus request/result
classes, not a new model.
An ack is not success. The server acks after validation and ingests asynchronously. Even an
invalid providerId is acked and then recorded as ERROR, which contradicts the proto comments. The only
confirmation is the request-status record, written after the data. So queryRequestStatus is now in
scope, with an await_request_statuses() poller.
Streaming semantics differ from the proto comments.ingestDataStream reports rejected ids on an error response, so the error result must keep them. A client-side error mid-stream gets no response,
but already-sent requests are still ingested. Separately, grpcio swallows exceptions raised by a
request iterator (verified); the wrappers re-raise the original.
Chunking is in scope: split_data_frame() for the server's ~4 MB message cap and 1-day bucket-span
cap, both of which a real pandas DataFrame hits immediately.
Scope — two PRs
PR A — client + unit tests.ingest_data(), ingest_data_stream(), iter_ingest_data_bidi_stream(), query_request_status() / RequestStatusQuery / await_request_statuses(), RegisterProviderApiResult.provider_id,
the non-scalar column builders, new data_frame() checks, and split_data_frame().
PR B — live tests + docs. Implement the skipped test_closed_loop_round_trip from #7. Move the two
integration tests that ingest through the raw stub onto the client. Add the missing live sample-status
filtering test. Add a cookbook ingestion recipe, and re-verify the query and sample-status recipes
against real data.
Out of scope
subscribeData() and subscribeDataEvent(): separate long-lived subscription lifecycle.
Add a Python client interface to the MLDP Ingestion Service (DpIngestionService) — the full ingestion surface, not just the unary call. The project currently has no real ingestion client outside the dp-desktop-app JavaFX demo, so this fills a genuine product gap: a facility standing up MLDP needs a straightforward Python path to get data in, not only to query it back out.
This is a peer sub-issue to the query interface (#7) and follows the same one-clean-sub-issue-per-PR pattern established by #5 (PV metadata) and #9 (machine config): params classes → _build_*_request() → _send_*() with three-tier error handling → *ApiResult, exposed via MldpClient (client.ingestion / the existing ingestion_client).
Why full-surface (not unary-only)
The bulk of the design work is the shared IngestDataRequest payload, which is identical across every ingestion RPC:
common.DataFrame = dataTimestamps (a DataTimestamps oneof: explicit timestampListor a samplingClock = startTime + periodNanos + count) + a large set of typed column arms (dataColumns generic DataValue, plus doubleColumns/int64Columns/stringColumns/imageColumns/structColumns/the array-column variants, etc.).
Response: IngestDataResponse with exceptionalResult | ackResult { numRows, numColumns }.
Because unary and streaming wrap the same payload + params model, unary-only would still require building 100% of that substrate and would risk a breaking reshape of the public params API when streaming lands later (the same "renaming later is breaking" argument that drove QueryParams/client.query in #7). Building the full surface once validates the shared types against every consumer up front.
Scope (phased internally; one PR)
Phase 1 — shared payload/params model + unary ingestData + unit tests
User-friendly params for a DataFrame: Python-native columns + a timestamps spec that accepts either an explicit timestamp list or a sampling clock (start + period + count), reusing to_timestamp().
ingestData() unary wrapper + IngestDataApiResult (surfacing ackResult.numRows/numColumns and the three-tier error handling).
This is the foundation, fully exercised by the unary path.
Phase 2 — streaming
ingestDataStream() (client-streaming / stream_unary) and ingestDataBidiStream() (stream_stream), built on the Phase 1 payload model.
Iterator-based ergonomics symmetric with the query client's streaming.
Phase 3 — closed-loop integration test + docs
Live ingest→query-back round-trip against the local ecosystem, asserting exact value round-trip, trimmed half-open range, dense column/timestamp alignment, and paging.
This also retroactively completes the query integration test in interface to v2 time-series data query API #7, which is intentionally deferred pending this client (a closed-loop query test needs a way to put known data in; assuming a test DB is pre-populated was explicitly rejected).
README + CLAUDE.md quick-start ("getting data in") section.
Notes / out of scope for now
queryRequestStatus, subscribeData — separate concerns; not part of the core ingest path. Can be filed separately if wanted.
Summary
Add the ingestion side of the Python client. Today
IngestionClientwraps onlyregisterProvider().Without an ingestion client, the closed-loop query test (#7), the cookbook's ingestion recipe, and the
live verification of sample-status filtering are all blocked.
What triage changed
data_frame.py/data_frame_conversions.pyas the substrate for this ticket (
ingestionDataFrameis the samecommon.DataFrameas calculations).So Phase 1 is a thin
IngestDataRequestParams(provider_id, frame, client_request_id)plus request/resultclasses, not a new model.
invalid
providerIdis acked and then recorded as ERROR, which contradicts the proto comments. The onlyconfirmation is the request-status record, written after the data. So
queryRequestStatusis now inscope, with an
await_request_statuses()poller.ingestDataStreamreports rejected ids on anerror response, so the error result must keep them. A client-side error mid-stream gets no response,
but already-sent requests are still ingested. Separately, grpcio swallows exceptions raised by a
request iterator (verified); the wrappers re-raise the original.
the matching checks in
data_frame(), which today lets hand-built columns through that the serverrejects.
querySamplesis scalar-only, so these can be verified live only up to ingestion until interface to v2 bucket-oriented query API (queryBuckets / queryBucketsStream) #16.split_data_frame()for the server's ~4 MB message cap and 1-day bucket-spancap, both of which a real pandas DataFrame hits immediately.
Scope — two PRs
PR A — client + unit tests.
ingest_data(),ingest_data_stream(),iter_ingest_data_bidi_stream(),query_request_status()/RequestStatusQuery/await_request_statuses(),RegisterProviderApiResult.provider_id,the non-scalar column builders, new
data_frame()checks, andsplit_data_frame().PR B — live tests + docs. Implement the skipped
test_closed_loop_round_tripfrom #7. Move the twointegration tests that ingest through the raw stub onto the client. Add the missing live sample-status
filtering test. Add a cookbook ingestion recipe, and re-verify the query and sample-status recipes
against real data.
Out of scope
subscribeData()andsubscribeDataEvent(): separate long-lived subscription lifecycle.queryBuckets).client.ingestionalias; the client staysclient.ingestion_client.Part of epic #10.
Original AI-drafted description
Summary
Add a Python client interface to the MLDP Ingestion Service (
DpIngestionService) — the full ingestion surface, not just the unary call. The project currently has no real ingestion client outside thedp-desktop-appJavaFX demo, so this fills a genuine product gap: a facility standing up MLDP needs a straightforward Python path to get data in, not only to query it back out.This is a peer sub-issue to the query interface (#7) and follows the same one-clean-sub-issue-per-PR pattern established by #5 (PV metadata) and #9 (machine config): params classes →
_build_*_request()→_send_*()with three-tier error handling →*ApiResult, exposed viaMldpClient(client.ingestion/ the existingingestion_client).Why full-surface (not unary-only)
The bulk of the design work is the shared
IngestDataRequestpayload, which is identical across every ingestion RPC:IngestDataRequest { providerId, clientRequestId, ingestionDataFrame: common.DataFrame }common.DataFrame=dataTimestamps(aDataTimestampsoneof: explicittimestampListor asamplingClock= startTime + periodNanos + count) + a large set of typed column arms (dataColumnsgenericDataValue, plusdoubleColumns/int64Columns/stringColumns/imageColumns/structColumns/the array-column variants, etc.).IngestDataResponsewithexceptionalResult|ackResult { numRows, numColumns }.Because unary and streaming wrap the same payload + params model, unary-only would still require building 100% of that substrate and would risk a breaking reshape of the public params API when streaming lands later (the same "renaming later is breaking" argument that drove
QueryParams/client.queryin #7). Building the full surface once validates the shared types against every consumer up front.Scope (phased internally; one PR)
Phase 1 — shared payload/params model + unary
ingestData+ unit testsDataFrame: Python-native columns + a timestamps spec that accepts either an explicit timestamp list or a sampling clock (start + period + count), reusingto_timestamp().ingestData()unary wrapper +IngestDataApiResult(surfacingackResult.numRows/numColumnsand the three-tier error handling).Phase 2 — streaming
ingestDataStream()(client-streaming /stream_unary) andingestDataBidiStream()(stream_stream), built on the Phase 1 payload model.Phase 3 — closed-loop integration test + docs
Notes / out of scope for now
queryRequestStatus,subscribeData— separate concerns; not part of the core ingest path. Can be filed separately if wanted.serializedDataColumnson the ingest side mirror the query-side deferral (interface to v2 time-series data query API #7 Q3) — default to the dense typed columns.Dependency
Blocks the query integration test in #7 (#7). Part of epic #10.