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
5 changes: 4 additions & 1 deletion CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -244,6 +244,8 @@ plan documents one change, `CLAUDE.md` documents the invariant it established.
- `src/dp_python_lib/client/data_frame.py` - Builders for `common.DataFrame`, the shared time-series payload (also ingestion's `ingestionDataFrame`, so #17 extends this rather than forking it): the `sampling_clock()` / `timestamp_list()` / `timestamp_count()` axis helpers **relocated here from `sample_status_client`** (and re-exported from it, so existing imports keep working), the typed scalar column builders (`double_column`, `float_column`, `int64_column`, `int32_column`, `bool_column`, `string_column`, `enum_column`), the legacy `data_column()` escape hatch (a `None` entry becomes an unset oneof — the only way to express a gap on a shared axis), the provenance helpers (`column_metadata`, `provenance`, `pv_source`, `calculations_source`), and `data_frame()` assembly, which routes columns by type and enforces the server's **shape** rules client-side (non-blank names, non-empty values, count match, name uniqueness across all types) while leaving its size caps server-side. Column names must be **non-blank**, not merely non-empty, and `timestamp_count()` validates a hand-built axis the way the builders do: a `SamplingClock` needs a positive `periodNanos` as well as a non-zero count, and a `TimestampList` must be **strictly increasing** (duplicates included — two samples cannot claim one instant). The read path rejects all of these, and anything `data_frame()` accepts must be readable back. Shared with `SampleStatusFrame`, which validates the same way. Sample counts come from the right field per kind: `dataValues` for a `DataColumn`, `images` for an `ImageColumn` (which has no `values` field at all), and `len(values) / prod(dims)` for an array column. An array column's sample count is `len(values) / prod(dims)`, and a value count that is not a **whole multiple** of `prod(dims)` is rejected rather than floor-divided into a passing count — the read path applies the same rule, so a frame this accepts is always one `data_frame_conversions` can read back. `data_column()` maps integers and floats by `numbers.Integral` / `numbers.Real` (and NumPy's bool by type name, without importing NumPy), so NumPy scalars map like their Python counterparts: `np.float64` subclasses `float` but `np.int64` and `np.bool_` subclass nothing here, and matching on exact Python type would accept some and reject others. #17 added the non-scalar builders: `double_array_column` / `float_array_column` / `int32_array_column` / `int64_array_column` / `bool_array_column` (samples as nested sequences or NumPy arrays, duck-typed; shape inferred and required identical across samples; flattened **row-major**, which is what the read path assumes; 1–3 dims; explicit `dims=` may shape flat samples or restate a shape, never reinterpret one), `image_column` (one descriptor per column, as in the proto), `struct_column`, and `serialized_column`. **Each kind's structural fields are checked in `_check_column()`, not only in the builders** — an enum's `enumId`, 1–3 array dims, an image's descriptor (positive width/height/channels, non-blank encoding), a struct's `schemaId`, a serialized column's `encoding` — because a hand-built column bypasses every builder and the server rejects all of these (the recurring #6 defect shape). `validate_data_frame(frame)` applies the same checks to an assembled frame (ingestion re-validates with it; its `caller=` keyword names the entry point an error message starts with, so a frame rejected by `IngestDataRequestParams` or `split_data_frame()` is not reported as a `data_frame()` error); `data_frame()` shares the column-list check rather than calling it, so messages keep the caller's index. `split_data_frame(frame, *, max_rows, max_bytes, max_span_nanos)` chunks lazily along the time axis: `max_bytes` bounds the **whole `IngestDataRequest`**, with room for two worst-case ids (`MAX_BUDGETED_ID_CHARS` = 256 chars × 4 UTF-8 bytes, computed from the message, not hard-coded); `max_span_nanos` is first-to-last timestamp, the server's bucket-span measure; `SamplingClock` chunk starts are integer nanoseconds. No limit defaults: `SERVER_DEFAULT_MAX_MESSAGE_BYTES` (4,096,000) is exported as a reference only, since the server value is deployment configuration. A frame with a `SerializedDataColumn` cannot be split
- `src/dp_python_lib/client/data_frame_conversions.py` - Reading a `DataFrame` back. Pure Python (no extras): `data_frame_timestamps()` (integer-nanosecond axis expansion), `column_values()` (standalone per-column converter, written so the bucket query #16 can reuse it; it yields exactly one entry per sample for every column kind, with the structural fields a payload cannot be interpreted without kept in companion accessors: array dims via `column_dimensions()` / `data_frame_column_dimensions()` so `[2,2]` and `[4]` stay distinguishable, an image's width/height/channels/encoding via `image_descriptor_dict()` / `data_frame_image_descriptors()`, and a struct's `schemaId` via `column_schema_id()` / `data_frame_schema_ids()`), `data_frame_columns()`, `column_metadata_dict()`. An axis that is set but empty is rejected rather than converted to a zero-row table. Behind `[analysis]`: `data_frame_to_pandas()` (UTC index built from int64 nanos, `ColumnMetadata` in `df.attrs`; each Series carries the **narrow dtype its column type implies** — `float32`/`int32`, not pandas' widened inference), `data_frame_from_pandas()` (dtype→typed column; a NaN anywhere is fail-loud, since a dense typed column cannot express a gap; rebuilds each column's `ColumnMetadata` from `df.attrs` via `column_metadata_from_dict()`). **A frame round-trips through pandas to byte equality** apart from the deliberate `SamplingClock`→`TimestampList` axis change: column types and provenance both survive. An `EnumColumn`'s codes are int32 and indistinguishable from a plain `Int32Column` by dtype, so its `enumId` rides in `df.attrs["enum_ids"]` — carried even under `exclude_column_metadata=True`, since without it the column cannot be rebuilt as an enum at all, and the `calculations_to_dataframes()` / `calculations_from_dataframes()` bridges. The pandas direction always emits a `TimestampList`, never an inferred `SamplingClock`; reads each instant as `Timestamp.value` (**always nanoseconds**, unlike a raw int64 view, which is in the index's own storage unit — a `datetime64[us]` index viewed as int64 lands 1000× too early); rejects `NaT`, whose integer form is a valid-looking instant; and rejects duplicate column names up front (a duplicated label makes `df[name]` a DataFrame, which would otherwise fail deep inside as pandas' "truth value of a Series is ambiguous"). `column_metadata_dict()` reports **only the origin arm actually set** on each provenance source — a PV source has no `calculations_column` key and vice versa — plus `time_range` as epoch nanoseconds when present; an absent range has no key rather than a fabricated `(0, 0)`. A pandas round trip preserves every value, dtype, and timestamp but **not column order**: a `DataFrame` stores each column kind in its own repeated field, so columns come back grouped by type
- `src/dp_python_lib/client/query_client.py` - v2 time-series query client (sample-oriented) exposed as `client.query`. Low-level wrappers `query_samples()` (unary, one resumable page) and `iter_query_samples()` (transparent paging), plus `iter_query_samples_stream()` (server-streaming, fire-and-consume, lazy). Queries are described by a kind-neutral `QueryParams` built from the `PvQuery` (`PV`) and `ConfigQuery` (`CFG`) criterion helpers; shares a `_build_query_spec()` seam so a future bucket request builder reuses it. Results wrap the raw `ColumnTable` (`.column_table`, `.next_page_token`); `.to_dataframe()`/`.to_numpy()` delegate to `query_conversions` (Phase 2, optional `[analysis]` extra)
- Bucket-oriented (#16 PR A; `plan/tickets/16/plan.md`): `query_buckets()` / `iter_query_buckets()` / `iter_query_buckets_stream()`, the exact contracts of the three sample methods over the same `QueryParams`. `_build_query_buckets_request()` **refuses** a `sample_status_filter` with a `ValueError` before any RPC (the server rejects its presence on buckets; dropping it would return unfiltered data to a caller who thinks it filtered), and forces `useSerializedColumns = False` (inert on buckets server-side) while passing `excludeColumnMetadata` through (honored on buckets, unlike samples). `QueryBucketsApiResult` exposes `.data_buckets` (a list, empty on error) and `.next_page_token`, and has `.to_dataframes()` but deliberately **no** `.to_dataframe()`: a bucket page is not one table. Both stream senders go through the private `_iter_stream()`, which keeps the yield-errors contract and error texts of the old sample-stream sender; it is still not `_dispatch`. `limit` counts **buckets** on this path, and a page may end early (with a token) at the server's byte budget
- `src/dp_python_lib/client/bucket_conversions.py` - Reading bucket query results (#16). Pure Python: `bucket_column()` (whichever arm is set, `SerializedDataColumn` included -- the only way to reach a serialized payload), `bucket_to_data_frame()` (a bucket as a one-column `common.DataFrame`, so every `data_frame_conversions` reader applies), `bucket_timestamps()`, `bucket_values()`, `trim_bucket()`, `buckets_by_pv()`. Behind `[analysis]`: `buckets_to_dataframes()` (`dict[pv, DataFrame]`, one column named for the PV) and `query_buckets_to_dataframes()` (unary paging with `max_buckets`; deliberately no streaming counterpart, since a PV spans messages). **Read-side checks only**: nothing here calls `validate_data_frame()` or `timestamp_count()`, because ingestion stores a *non-decreasing* `TimestampList` with repeats allowed and the write path's strictly-increasing rule would make such a bucket unreadable. `trim_bucket()` is exact and half-open, reuses `data_frame._slice_frame()` (a `SamplingClock` stays a clock with a shifted start), finds both bounds with `bisect_left`, and runs the read path's alignment check first -- `_slice_frame()` silently slices an array column with absent dims in blocks of 1. A decreasing axis, a serialized bucket, and `begin >= end` all raise. Grouping does not rely on the server's (pvName, firstTime) order, and overlapping buckets are kept, not deduplicated. A PV whose buckets differ in kind or structure (`enumId`, dims, image descriptor, `schemaId`) raises. `df.attrs` is **rebuilt** after `pd.concat`, never inherited: `buckets` (per-bucket first/last nanos, count, provider, stored `column_name`, and metadata -- `None`, not an empty summary, when the bucket carries none, e.g. after a query that excluded it), `column_metadata` only when every bucket agrees, and the structural `enum_ids` / `dimensions` / `image_descriptors` / `schema_ids`, kept even under `exclude_column_metadata`
- `src/dp_python_lib/client/query_conversions.py` - Pythonic conversions for query results (optional `[analysis]` extra: pandas/numpy/openpyxl, imported lazily). `data_value_to_python()` (oneof extractor: scalars→native, timestamp→epoch-nanos, array→list, structure→dict, image→`Image` wrapper, fail-loud on unhandled arm), `column_table_to_dataframe()` (UTC datetime index + one column per DataColumn; dense-alignment and duplicate-column-name fail-loud; ColumnMetadata in `df.attrs`), `column_table_to_numpy()` (dict of 1-D arrays; complex arms stay 1-D object arrays rather than collapsing to 2-D), `dataframe_to_excel()` (thin `to_excel()` wrapper: row-limit guard, tz-drop, complex-cell stringification), and `query_samples_to_dataframe()`/`stream_query_samples_to_dataframes()` whole-query conveniences (unary concats by column name; streaming yields per-page frames lazily)
- `src/dp_python_lib/client/service_api_client_base.py` - Base class for the service clients: owns the channel and the one-per-client gRPC stub, and provides `_dispatch()`, the shared three-tier sender that all 18 unary `_send_*` methods delegate to
- `src/dp_python_lib/client/query_support.py` - Helpers shared by the criteria-based query clients, currently `check_at_most_one_text_criterion()` (Mongo cannot AND two `$text` clauses). It is generic over the different criterion types because each names its oneof `criterion` and its full-text arm `textCriterion`. New shared query helpers belong here rather than in whichever feature client happened to need one first — importing a private name across feature modules makes the importing module's dependencies misleading
Expand All @@ -267,7 +269,8 @@ plan documents one change, `CLAUDE.md` documents the invariant it established.
- `tests/unit/test_time_conversions.py` - Unit tests for the shared time converters (`to_timestamp()` input forms and the naive-datetime/bool/unsupported-type rejections; `to_epoch_nanos()` exactness and its round trip with `to_timestamp()`; `from_epoch_nanos()` and its round trip)
- `tests/unit/test_data_frame.py` - Unit tests for the data_frame builders (axis relocation, each typed column, `data_column()` bool-before-int and unset-oneof handling, provenance helpers, and every `data_frame()` shape rule incl. array dims and serialized-column name-only checks)
- `tests/unit/test_data_frame_conversions.py` - Unit tests for data_frame_conversions (nanosecond-exact expansion, per-column conversion incl. array reshaping, duplicate-name fail-loud, and — skipping cleanly without `[analysis]` — the pandas round trip, dtype mapping, and NaN fail-loud)
- `tests/unit/test_query_client.py` - Unit tests for QueryClient (request building, three-tier error handling, unary paging, streaming, `PvQuery`/`ConfigQuery` helpers, `QueryParams` validation)
- `tests/unit/test_query_client.py` - Unit tests for QueryClient (request building, three-tier error handling, unary paging, streaming, `PvQuery`/`ConfigQuery` helpers, `QueryParams` validation; for buckets, the `sample_status_filter` refusal before any RPC, forced `useSerializedColumns`, and the error texts of the shared `_iter_stream()`)
- `tests/unit/test_bucket_conversions.py` - Unit tests for bucket_conversions (every column arm through the frame view and values, nanosecond-exact axes, trimming at both bounds incl. repeated timestamps at each bound, clock start shift, decreasing-axis / absent-dims / serialized / `begin >= end` refusals, grouping, kind and structure mismatches, per-bucket vs uniform metadata in `attrs`, `max_buckets`; pandas tests skip cleanly without `[analysis]`)
- `tests/unit/test_query_conversions.py` - Unit tests for query_conversions (each DataValue arm, dense-alignment and duplicate-column-name fail-loud, int-gap float-upcast, timestamp columns, 1-D object arrays for complex arms, metadata in attrs, concat-by-name, Excel row-limit/stringification/native-bytes; DataFrame/NumPy/Excel tests skip cleanly when the `[analysis]` extra is absent)
- `pyproject.toml` - Project metadata and dependencies
- Generated gRPC stubs include services for:
Expand Down
44 changes: 44 additions & 0 deletions doc/release-notes/NEXT.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,7 @@ person cutting the release has any reason to re-read.
- [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)
- [Querying whole buckets (#16)](#querying-whole-buckets-issue-16)
- [Cutting the release](#cutting-the-release)

---
Expand Down Expand Up @@ -217,6 +218,49 @@ Re-running the older recipes against real data also corrected two things they cl
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).

## Querying whole buckets (Issue #16)

`client.query` now wraps the bucket-oriented v2 query, `queryBuckets` / `queryBucketsStream`. Where
`query_samples()` returns an aligned, trimmed table of scalar columns, a bucket query returns the
archive's stored units: each `DataBucket` holds one PV's column, in the type it was ingested as, over
that bucket's own time axis. It is the way to read back array, image, struct, and serialized columns,
and the first query path that returns the column metadata (provenance included) stored with the data.

- **`query_buckets()`, `iter_query_buckets()`, and `iter_query_buckets_stream()`** take the same
`QueryParams` as the sample methods and behave the same way: one page, every page (raising
`RuntimeError` on a page error), or the server stream (raising on a mid-stream error).
`QueryBucketsApiResult` exposes `data_buckets` and `next_page_token`.
- **A new module, `bucket_conversions`**, reads them. In plain Python: `bucket_column()`,
`bucket_values()`, `bucket_timestamps()` (integer nanoseconds), `bucket_to_data_frame()` (a bucket
as a one-column `common.DataFrame`, so the existing `data_frame_conversions` readers apply), and
`buckets_by_pv()`. With the `[analysis]` extra: `buckets_to_dataframes()`, which returns one pandas
DataFrame per PV, and `query_buckets_to_dataframes()`, which runs the whole query first (with a
`max_buckets` cap). `QueryBucketsApiResult.to_dataframes()` converts one page.

Behaviors worth knowing:

- **Buckets come back whole.** The server returns every bucket that overlaps `[begin, end)`
without trimming it, so a bucket at either edge carries samples outside the range. Trimming is
opt-in: `trim_bucket()`, `time_range=` on the conversions, or `trim=True` on
`query_buckets_to_dataframes()`. It is exact and half-open, like `query_samples()`. It does not
remove samples that fall in a gap between configuration intervals when `config_criteria` is used;
use `query_samples()` when that matters.
- **`limit` counts buckets on this path, not rows**, and a page can also end early, with a page
token, once it reaches the server's message-size budget. The server silently caps the limit at its
configured maximum (100,000 by default).
- **A `QueryParams` with a `sample_status_filter` is refused** with a `ValueError` before any call
is made. The server rejects status filtering on bucket queries, and dropping the filter silently
would return data the caller believed was filtered.
- **A serialized column is passed through, not decoded.** It comes back only if it was ingested
that way. `bucket_column()` returns it with its `encoding` and `payload`; the value readers and
pandas conversions raise rather than drop it.
- **A PV's buckets are kept as stored.** Overlapping buckets from separate ingests are all kept,
so a per-PV frame's index can repeat instants. A PV whose buckets differ in column type or
structure (say, ingested as double and later as int32) raises, naming both buckets. Per-bucket
provider and column metadata is in `df.attrs["buckets"]`.

See [#16](https://github.com/osprey-dcs/dp-python-lib/issues/16) and `plan/tickets/16/plan.md`.

## Installing

```bash
Expand Down
2 changes: 2 additions & 0 deletions src/dp_python_lib/client/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -100,6 +100,7 @@
from dp_python_lib.client.query_client import (
ConfigQuery,
PvQuery,
QueryBucketsApiResult,
QueryClient,
QueryParams,
QuerySamplesApiResult,
Expand Down Expand Up @@ -157,6 +158,7 @@
"PvMetadataQuery",
"PvQuery",
"QueryAnnotationsApiResult",
"QueryBucketsApiResult",
"QueryClient",
"QueryConfigurationActivationsApiResult",
"QueryConfigurationsApiResult",
Expand Down
Loading
Loading