Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
113 changes: 106 additions & 7 deletions CLAUDE.md

Large diffs are not rendered by default.

13 changes: 7 additions & 6 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -88,8 +88,13 @@ for their service.
bridges to pandas under the optional `[analysis]` extra.
- **Export** — `client.annotation.export`. Export a saved dataset, ad-hoc data blocks, and/or
calculations to HDF5, CSV, or XLSX. The file is written on the server; there is no retrieval RPC.
- **Provider registration** — `client.ingestion_client.register_provider()`. The rest of the
ingestion API is not yet implemented.
- **Ingestion** — `client.ingestion_client`. Register a provider with `register_provider()`, then
ingest a `data_frame` frame with `ingest_data()`, `ingest_data_stream()` (many requests, one
summary), or `iter_ingest_data_bidi_stream()` (an ack or reject per request). An ack means only
that a request passed validation: confirm it landed with `await_request_statuses()` or
`query_request_status()`. `data_frame.split_data_frame()` cuts a large frame into chunks under
the server's message-size and time-span limits, and the `data_frame` builders now cover array,
image, struct, and serialized columns as well as scalars.

**Supporting framework:** YAML + environment-variable configuration (`MLDP_*`, via
pydantic-settings), TLS-capable channel creation, hierarchical logging, three-tier error handling
Expand All @@ -105,10 +110,6 @@ older than your `dp_python_lib` will not implement everything listed here. The
**Low-level API coverage**

- **Ingestion Service**
- `ingestData()` / `ingestDataStream()` / `ingestDataBidiStream()` — full ingestion client with a
shared DataFrame payload model ([issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17);
also unblocks the closed-loop query integration test and the cookbook's ingestion recipe)
- `queryRequestStatus()` — async status of ingestion requests
- `subscribeData()` — receive data for specified PVs from the ingestion stream
- **Query Service**
- `queryBuckets()` / `queryBucketsStream()` — raw data buckets
Expand Down
56 changes: 56 additions & 0 deletions doc/release-notes/NEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ person cutting the release has any reason to re-read.
- [Ready for typed gRPC stubs (#61)](#ready-for-typed-grpc-stubs-issue-61)
- [Environment variables override the config file (#19)](#environment-variables-override-the-config-file-issue-19)
- [Detecting open-ended activations (#26)](#detecting-open-ended-activations-issue-26)
- [Ingesting data (#17)](#ingesting-data-issue-17)
- [Cutting the release](#cutting-the-release)

---
Expand Down Expand Up @@ -140,6 +141,61 @@ not have its id.

See [#26](https://github.com/osprey-dcs/dp-python-lib/issues/26).

## Ingesting data (Issue #17)

The library can now put data into MLDP, not just read it. `client.ingestion_client`, which
previously offered only `register_provider()`, now covers data ingestion and request status
(`subscribeData()` is not wrapped yet):

- **`ingest_data()`** sends one request. **`ingest_data_stream()`** sends many on one call and
returns a single summary. **`iter_ingest_data_bidi_stream()`** yields each request's ack or
reject as it arrives. To stop a bidi stream early, close its iterator, most simply with
`contextlib.closing`: that cancels the call. A plain `break` does not, while the iterator is still
referenced, and gRPC keeps sending the remaining requests in the background.
- **`await_request_statuses()`** and **`query_request_status()`** report whether an ingestion
actually landed. This matters because **an ack means only that a request passed validation.**
The server queues the data and writes it afterwards, so a request can be acked and still fail.
An unknown provider id fails that way, and so does ingesting the same PV with the same first
timestamp twice. The request-status document, written once the data is stored, is the only
confirmation. Capture the time before sending and pass it as `since`.
- **`IngestDataRequestParams`** takes a `common.DataFrame` built with the existing `data_frame`
builders, or with `data_frame_from_pandas()`. Its request id defaults to a generated uuid, since
the server does not enforce unique ids.
- **`data_frame.split_data_frame()`** cuts a large frame into chunks that fit the server's inbound
message limit (4,096,000 bytes by default, about 500,000 doubles) and its one-day bucket span.
**`chunked_request_params()`** turns the chunks into requests with correlated ids, lazily, ready
for `ingest_data_stream()`. No limit is assumed: pass `SERVER_DEFAULT_MAX_MESSAGE_BYTES` to
target a default server.
- **New column builders** for the kinds that had none: `double_array_column()` and its float, int32,
int64, and bool siblings (samples as nested lists or NumPy arrays), `image_column()`,
`struct_column()`, and `serialized_column()`.
- **`RegisterProviderApiResult` gains `provider_id` and `is_new_provider`**, so the id no longer
has to be dug out of `.response.registrationResult`.
- **`from_epoch_nanos()`**, the inverse of `to_epoch_nanos()`, is exported.

Two behaviors worth knowing:

- **A rejected request inside `ingest_data_stream()` does not raise.** The result has `is_error`
set and lists `rejected_request_ids`; the other requests were accepted.
- **If your own request generator raises during a stream, you get your exception back**, not
gRPC's generic "Exception iterating requests!". Requests sent before it may or may not have
reached the server, so check their status.

**Behavior change: `data_frame()` rejects hand-built columns it used to pass.** If you build
column messages from the protobuf types directly and pass them to `data_frame()`, it now raises
`ValueError` for an `EnumColumn` without an `enumId`, an array column with more than three
dimensions, an `ImageColumn` without a complete image descriptor, a `StructColumn` without a
`schemaId`, and a `SerializedDataColumn` without an `encoding`. The server rejects all of these
anyway, so nothing that used to be accepted end to end is lost. `enum_column()` likewise now
rejects a whitespace-only `enum_id`.

Non-scalar columns (arrays, images, structs) can be ingested, but not yet read back through the
query API, which returns scalar columns only; reading them back is
[#16](https://github.com/osprey-dcs/dp-python-lib/issues/16).

See [#17](https://github.com/osprey-dcs/dp-python-lib/issues/17) and the Ingestion API section of
[`CLAUDE.md`](https://github.com/osprey-dcs/dp-python-lib/blob/main/CLAUDE.md).

## Installing

```bash
Expand Down
18 changes: 16 additions & 2 deletions plan/tickets/17/plan.md
Original file line number Diff line number Diff line change
Expand Up @@ -237,6 +237,18 @@ All dp-service citations are `origin/main` @ `7e8b2e6`, paths relative to
`RpcError`) instead of returning a gRPC error. Its message says how many requests had already been sent,
because those were ingested (T4). Unit-tested against an in-process grpcio server, since a mocked stub
would not reproduce grpcio's swallowing.
- *Corrected in implementation (2026-09-30).* Two details above were wrong. (1) "Already sent" does not
mean "ingested", or even "received": the in-process tests show grpcio's cancellation racing the sends, so
the server holds some prefix of the requests handed over, possibly none. The note therefore says "any that
reached the server are ingested -- check request status" (T10's "have already been ingested" is corrected
the same way). (2) The count cannot go in the exception's *message* without changing its type or mutating
its args, so it travels as a PEP 678 note (Python 3.11+) and is always logged.
- *Added in implementation.* Closing the iterator `iter_ingest_data_bidi_stream()` returns cancels the call.
Without that, grpcio keeps pulling and sending the caller's requests on its own thread after the caller has
walked away; an in-process test showed all 1,000 queued requests drained into the server. The cancel needs an
explicit close: a `break` out of a loop over a generator that is still referenced leaves it suspended until
garbage collection, so the documented way to stop early is `with contextlib.closing(...)` (corrected in PR
review, 2026-09-30; the first draft said abandoning the iterator was enough).

- **D6 — `queryRequestStatus` gets a criterion helper and a poller.**
- `RequestStatusQuery` (`RS`): `provider_id(id)`, `provider_name(name)`, `request_id(id)`,
Expand Down Expand Up @@ -337,8 +349,10 @@ All dp-service citations are `origin/main` @ `7e8b2e6`, paths relative to

**`src/dp_python_lib/client/data_frame.py`**
- Extract `validate_data_frame(frame)` from `data_frame()`: axis via `timestamp_count()`, at least one
column, then `_check_column()` on every column across all 16 arms. `data_frame()` calls it after
assembly.
column, then `_check_column()` on every column across all 16 arms. *(As implemented, `data_frame()` does
not call it after assembly: both share one column-list check, which `data_frame()` runs on the caller's list
before routing, so an unsupported type is caught before it must be routed and messages keep the caller's
index.)*
- Extend `_check_column()` with T7's rules: enum `enumId` non-blank; array dims count 1–3; image
descriptor present with positive width/height/channels and non-blank encoding; struct `schemaId`
non-blank; serialized `encoding` non-blank.
Expand Down
37 changes: 36 additions & 1 deletion src/dp_python_lib/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -16,19 +16,29 @@
# `import dp_python_lib.client.data_frame` would hand back the function instead of the module -- breaking the
# documented `from dp_python_lib.client import data_frame as dfb` usage. Reach it as dfb.data_frame(...).
from dp_python_lib.client.data_frame import (
bool_array_column,
bool_column,
calculations_source,
column_metadata,
data_column,
double_array_column,
double_column,
enum_column,
float_array_column,
float_column,
image_column,
int32_array_column,
int32_column,
int64_array_column,
int64_column,
provenance,
pv_source,
serialized_column,
split_data_frame,
string_column,
struct_column,
timestamp_count,
validate_data_frame,
)
from dp_python_lib.client.dataset_client import (
DataSetClient,
Expand All @@ -48,9 +58,16 @@
calculations_spec,
)
from dp_python_lib.client.ingestion_client import (
IngestDataApiResult,
IngestDataRequestParams,
IngestDataStreamApiResult,
IngestionClient,
IngestionRequestStatus,
QueryRequestStatusApiResult,
RegisterProviderApiResult,
RegisterProviderRequestParams,
RequestStatusQuery,
chunked_request_params,
)
from dp_python_lib.client.machine_config_client import (
ConfigurationActivationQuery,
Expand Down Expand Up @@ -101,7 +118,7 @@
timestamp_list,
)
from dp_python_lib.client.sample_status_conversions import SampleStatusRow
from dp_python_lib.client.time_conversions import TimestampInput, to_epoch_nanos, to_timestamp
from dp_python_lib.client.time_conversions import TimestampInput, from_epoch_nanos, to_epoch_nanos, to_timestamp

__all__ = [
"AnnotationClient",
Expand Down Expand Up @@ -129,7 +146,11 @@
"GetConfigurationApiResult",
"GetDataSetApiResult",
"GetPvMetadataApiResult",
"IngestDataApiResult",
"IngestDataRequestParams",
"IngestDataStreamApiResult",
"IngestionClient",
"IngestionRequestStatus",
"MachineConfigClient",
"MldpClient",
"PvMetadataClient",
Expand All @@ -142,11 +163,13 @@
"QueryDataSetsApiResult",
"QueryParams",
"QueryPvMetadataApiResult",
"QueryRequestStatusApiResult",
"QuerySampleStatusesApiResult",
"QuerySampleStatusesRequestParams",
"QuerySamplesApiResult",
"RegisterProviderApiResult",
"RegisterProviderRequestParams",
"RequestStatusQuery",
"SampleStatusClient",
"SampleStatusColumn",
"SampleStatusFilter",
Expand All @@ -167,24 +190,36 @@
"TimestampInput",
"activation_end_time",
"activation_is_open",
"bool_array_column",
"bool_column",
"calculations",
"calculations_source",
"calculations_spec",
"chunked_request_params",
"column_metadata",
"data_block",
"data_column",
"double_array_column",
"double_column",
"enum_column",
"float_array_column",
"float_column",
"from_epoch_nanos",
"image_column",
"int32_array_column",
"int32_column",
"int64_array_column",
"int64_column",
"provenance",
"pv_source",
"sampling_clock",
"serialized_column",
"split_data_frame",
"string_column",
"struct_column",
"timestamp_count",
"timestamp_list",
"to_epoch_nanos",
"to_timestamp",
"validate_data_frame",
]
Loading
Loading