From 0900904c5902bf53d3937020df8aef8b49c69d97 Mon Sep 17 00:00:00 2001 From: Craig McChesney Date: Thu, 1 Oct 2026 10:56:15 -0600 Subject: [PATCH 1/3] feat: bucket-oriented v2 query client and bucket_conversions (#16 PR A) Wrap queryBuckets / queryBucketsStream on QueryClient (query_buckets, iter_query_buckets, iter_query_buckets_stream) over the existing QueryParams, refusing sample_status_filter before any RPC. Both stream senders now share _iter_stream(). Add bucket_conversions: a bucket viewed as a one-column DataFrame, exact half-open trim_bucket() reusing _slice_frame() with read-side checks only, grouping by PV, and per-PV pandas frames with attrs rebuilt after concat. Refs osprey-dcs/dp-python-lib#16 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01CajTMkjkkzeoWXLSgMpn5k --- CLAUDE.md | 5 +- doc/release-notes/NEXT.md | 44 ++ src/dp_python_lib/client/__init__.py | 2 + .../client/bucket_conversions.py | 459 +++++++++++++++ .../client/data_frame_conversions.py | 12 +- src/dp_python_lib/client/query_client.py | 372 ++++++++++-- tests/unit/test_bucket_conversions.py | 542 ++++++++++++++++++ tests/unit/test_data_frame_conversions.py | 6 +- tests/unit/test_query_client.py | 265 ++++++++- 9 files changed, 1646 insertions(+), 61 deletions(-) create mode 100644 src/dp_python_lib/client/bucket_conversions.py create mode 100644 tests/unit/test_bucket_conversions.py diff --git a/CLAUDE.md b/CLAUDE.md index ad9e163..ba651ac 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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, metadata), `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 @@ -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: diff --git a/doc/release-notes/NEXT.md b/doc/release-notes/NEXT.md index a069136..df34872 100644 --- a/doc/release-notes/NEXT.md +++ b/doc/release-notes/NEXT.md @@ -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) --- @@ -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 diff --git a/src/dp_python_lib/client/__init__.py b/src/dp_python_lib/client/__init__.py index b2a487d..04e4978 100644 --- a/src/dp_python_lib/client/__init__.py +++ b/src/dp_python_lib/client/__init__.py @@ -100,6 +100,7 @@ from dp_python_lib.client.query_client import ( ConfigQuery, PvQuery, + QueryBucketsApiResult, QueryClient, QueryParams, QuerySamplesApiResult, @@ -157,6 +158,7 @@ "PvMetadataQuery", "PvQuery", "QueryAnnotationsApiResult", + "QueryBucketsApiResult", "QueryClient", "QueryConfigurationActivationsApiResult", "QueryConfigurationsApiResult", diff --git a/src/dp_python_lib/client/bucket_conversions.py b/src/dp_python_lib/client/bucket_conversions.py new file mode 100644 index 0000000..0132c5a --- /dev/null +++ b/src/dp_python_lib/client/bucket_conversions.py @@ -0,0 +1,459 @@ +""" +Pythonic conversions for bucket query results (issue #16; plan/tickets/16/plan.md). + +A bucket query (QueryClient.query_buckets() and friends) returns the archive's stored units whole: each DataBucket +holds ONE PV's column, in its stored typed form, over that bucket's own time axis. This module reads them. + +The pure-Python half has no third-party dependencies. It views a bucket as a one-column common.DataFrame, so every +data_frame_conversions reader applies unchanged, and trims a bucket to [begin, end) exactly at the proto level. +pandas is imported lazily inside the entry points that need it ([analysis] extra), so importing this module never +requires the extra. + +Design decisions (see the plan's D4-D8): + - Buckets come back UNTRIMMED: the server selects every bucket overlapping [begin, end) and returns it whole. + Trimming is opt-in (trim_bucket(), time_range=), exact, and half-open. It cannot remove samples in a gap + between configuration intervals; use querySamples() for that. + - Assembly groups by pvName and orders each PV's buckets by first timestamp itself, rather than relying on the + server's (undocumented) arrival order. Overlapping or duplicate timestamps across buckets are kept. + - Read-side checks only. A bucket is server data, and ingestion allows a non-decreasing TimestampList with + repeated timestamps, which this library's write path (validate_data_frame(), timestamp_count()) would refuse. + Nothing here calls those. + - A SerializedDataColumn bucket is passed through (bucket_column(), bucket_to_data_frame()) but refused wherever + its values would have to be interpreted, since its payload is opaque and the encoding contract is the caller's. +""" + +from bisect import bisect_left +from collections.abc import Iterable +from typing import Any + +from dp_python_lib.client import data_frame_conversions as dfc + +# _slice_frame() and _time_range() are private to data_frame.py, which is the shared, kind-neutral home of the +# DataFrame builders rather than a feature module. Reusing them keeps one implementation of the integer-nanosecond +# SamplingClock shift and of the "begin strictly before end" rule and its message. +from dp_python_lib.client.data_frame import _COLUMN_FIELD_BY_TYPE, _slice_frame, _time_range +from dp_python_lib.client.sample_status_conversions import expand_data_timestamps +from dp_python_lib.client.time_conversions import TimestampInput, to_epoch_nanos +from dp_python_lib.grpc import common_pb2 + + +def _require_pandas(): + """Imports and returns pandas, or raises an actionable error if the optional [analysis] extra is missing.""" + try: + import pandas + except ImportError as e: + raise ImportError( + "pandas is required for bucket DataFrame conversions. Install the optional analysis extra: " + 'pip install "dp-python-lib[analysis]"' + ) from e + return pandas + + +# ---------------------------------------------------------------------- +# one bucket +# ---------------------------------------------------------------------- + + +def bucket_column(bucket: common_pb2.DataBucket) -> Any: + """ + Returns a bucket's column message, whichever kind it is -- a SerializedDataColumn included. + + This is the way to reach a serialized bucket's payload and encoding, which the value readers refuse. + + :param bucket: The DataBucket. + :return: The column message set in bucket.dataValues (DoubleColumn, ImageColumn, SerializedDataColumn, ...). + :raises ValueError: if the bucket carries no column. + """ + arm = bucket.dataValues.WhichOneof("values") if bucket.HasField("dataValues") else None + if arm is None: + raise ValueError(f"DataBucket for PV '{bucket.pvName}' carries no column (dataValues is unset)") + return getattr(bucket.dataValues, arm) + + +def _require_axis(bucket: common_pb2.DataBucket) -> None: + """Rejects a bucket whose time axis is unset, naming its PV.""" + if bucket.dataTimestamps.WhichOneof("value") is None: + raise ValueError(f"DataBucket for PV '{bucket.pvName}' has no dataTimestamps; every bucket carries a time axis") + + +def bucket_to_data_frame(bucket: common_pb2.DataBucket) -> common_pb2.DataFrame: + """ + Views a bucket as a one-column common.DataFrame over the bucket's own time axis. + + Every data_frame_conversions reader then applies unchanged: data_frame_columns(), data_frame_to_pandas(), and + the companion accessors for array dims, image descriptors, and struct schema ids. A serialized column lands in + serializedDataColumns, where those readers skip it by design. + + The frame is not validated with validate_data_frame(): that is the write path's check, and it rejects the + repeated timestamps a stored bucket may legitimately carry. + + :param bucket: The DataBucket. + :return: A DataFrame carrying the bucket's axis and its single column. + :raises ValueError: if the bucket has no time axis or no column. + """ + _require_axis(bucket) + column = bucket_column(bucket) + + frame = common_pb2.DataFrame() + frame.dataTimestamps.CopyFrom(bucket.dataTimestamps) + getattr(frame, _COLUMN_FIELD_BY_TYPE[type(column)]).add().CopyFrom(column) + return frame + + +def bucket_timestamps(bucket: common_pb2.DataBucket) -> list[int]: + """ + Expands a bucket's time axis into one integer epoch-nanosecond value per sample. + + A SamplingClock is expanded in integer nanoseconds, so the result is exact at present-day epochs. + + :param bucket: The DataBucket. + :return: Epoch nanoseconds for each sample, in axis order. + :raises ValueError: if the axis is unset, empty, or malformed. + """ + _require_axis(bucket) + epoch_nanos = expand_data_timestamps(bucket.dataTimestamps) + if not epoch_nanos: + raise ValueError(f"DataBucket for PV '{bucket.pvName}' has a time axis describing no timestamps") + return epoch_nanos + + +def _refuse_serialized(bucket: common_pb2.DataBucket, column: Any, what: str) -> None: + """Raises if column is a SerializedDataColumn, naming the PV, its encoding, and the way to reach it.""" + if isinstance(column, common_pb2.SerializedDataColumn): + raise ValueError( + f"DataBucket for PV '{bucket.pvName}' holds a SerializedDataColumn (encoding {column.encoding!r}), " + f"whose payload is opaque, so it cannot be {what}; read it with bucket_column() and decode it per its " + f"encoding" + ) + + +def _aligned_values(bucket: common_pb2.DataBucket, column: Any, n_rows: int) -> list: + """ + Converts a bucket's column and checks it holds exactly one sample per timestamp. + + This is the read path's count rule (data_frame_conversions.column_values()), so an array column whose dims are + missing or do not divide its values raises here rather than being sliced in the wrong-sized blocks. + + :param bucket: The DataBucket, named in the error. + :param column: Its column, already known not to be serialized. + :param n_rows: The number of timestamps on its axis. + :return: One value per sample. + :raises ValueError: if the column cannot be converted or its sample count differs from n_rows. + """ + values = dfc.column_values(column) + if len(values) != n_rows: + raise ValueError( + f"DataBucket for PV '{bucket.pvName}' has {len(values)} samples but its time axis has {n_rows} " + f"timestamps; a bucket's column must be index-aligned with its axis" + ) + return values + + +def bucket_values(bucket: common_pb2.DataBucket) -> list: + """ + Extracts one Python value per sample from a bucket's column, in axis order. + + The mapping is data_frame_conversions.column_values()'s: native scalars, integer codes for an enum (its enumId is + on bucket_column()), one list per sample for an array (dims via data_frame_conversions.column_dimensions()), one + bytes payload per sample for an image, and None for a gap in a legacy DataColumn. + + :param bucket: The DataBucket. + :return: One value per sample, parallel to bucket_timestamps(). + :raises ValueError: if the bucket is malformed or not aligned with its axis, or holds a SerializedDataColumn. + """ + column = bucket_column(bucket) + _refuse_serialized(bucket, column, "converted to values") + return _aligned_values(bucket, column, len(bucket_timestamps(bucket))) + + +def _range_nanos(begin: TimestampInput, end: TimestampInput) -> tuple[int, int]: + """ + Converts a (begin, end) pair to epoch nanoseconds, requiring begin strictly before end. + + A reversed or empty range is rejected rather than trimming to nothing, which would be indistinguishable from a + valid range that simply holds no samples. + + :raises ValueError: if begin is not strictly before end, or either bound is not a supported time input. + """ + time_range = _time_range((begin, end)) + return to_epoch_nanos(time_range.beginTime), to_epoch_nanos(time_range.endTime) + + +def _first_out_of_order(epoch_nanos: list[int]) -> int | None: + """Returns the first position whose timestamp is earlier than its predecessor's, or None if non-decreasing.""" + for position in range(1, len(epoch_nanos)): + if epoch_nanos[position] < epoch_nanos[position - 1]: + return position + return None + + +def _trim_nanos(bucket: common_pb2.DataBucket, begin_nanos: int, end_nanos: int) -> common_pb2.DataBucket | None: + """trim_bucket() over an already-validated range in epoch nanoseconds.""" + column = bucket_column(bucket) + _refuse_serialized(bucket, column, "trimmed") + + epoch_nanos = bucket_timestamps(bucket) + # The read-side alignment check. _slice_frame() assumes a validated frame, and on an array column with no dims + # it would fall back to one value per row and slice silently wrong. + _aligned_values(bucket, column, len(epoch_nanos)) + + # A stored axis is non-decreasing (ingestion enforces it, repeats allowed), so the rows in [begin, end) are one + # contiguous span and bisect_left finds both of its ends -- with repeats too: every copy of a timestamp equal to + # begin is kept, and every copy equal to end is dropped. A decreasing axis has no such span. + out_of_order = _first_out_of_order(epoch_nanos) + if out_of_order is not None: + raise ValueError( + f"DataBucket for PV '{bucket.pvName}' has a decreasing time axis at position {out_of_order} " + f"({epoch_nanos[out_of_order]} ns follows {epoch_nanos[out_of_order - 1]} ns), so it cannot be trimmed " + f"to a contiguous span" + ) + + start = bisect_left(epoch_nanos, begin_nanos) + end = bisect_left(epoch_nanos, end_nanos) + if start >= end: + return None + + sliced = _slice_frame(bucket_to_data_frame(bucket), start, end) + trimmed = common_pb2.DataBucket() + trimmed.CopyFrom(bucket) + trimmed.dataTimestamps.CopyFrom(sliced.dataTimestamps) + # The copy has the same column arm set, so bucket_column() returns it in place, ready to overwrite. + bucket_column(trimmed).CopyFrom(getattr(sliced, _COLUMN_FIELD_BY_TYPE[type(column)])[0]) + return trimmed + + +def trim_bucket( + bucket: common_pb2.DataBucket, begin: TimestampInput, end: TimestampInput +) -> common_pb2.DataBucket | None: + """ + Returns a copy of a bucket holding only its samples in the half-open range [begin, end). + + Exact: timestamps are compared as integer epoch nanoseconds, so a sample exactly at begin is kept and one exactly + at end is dropped, as on the querySamples() path. A SamplingClock bucket stays a SamplingClock with its start + shifted in integer nanoseconds. The PV, provider, and column metadata are carried over unchanged. + + This trims to [begin, end) only. With a configuration selector the server can return a bucket spanning a gap + between activation intervals, and trimming does not remove the samples in that gap. + + :param bucket: The DataBucket to trim. It is not modified. + :param begin: Inclusive start (tz-aware datetime, epoch seconds, or common.Timestamp). + :param end: Exclusive end, strictly after begin. + :return: The trimmed copy, or None when no sample falls in [begin, end). + :raises ValueError: if begin is not strictly before end; if the bucket is malformed, holds a SerializedDataColumn + (an opaque payload cannot be sliced), is not aligned with its axis, or has a decreasing time axis. + """ + begin_nanos, end_nanos = _range_nanos(begin, end) + return _trim_nanos(bucket, begin_nanos, end_nanos) + + +# ---------------------------------------------------------------------- +# many buckets +# ---------------------------------------------------------------------- + + +def _first_nanos(bucket: common_pb2.DataBucket) -> int: + """A bucket's first timestamp, read without expanding the axis.""" + _require_axis(bucket) + axis = bucket.dataTimestamps + if axis.WhichOneof("value") == "samplingClock": + return to_epoch_nanos(axis.samplingClock.startTime) + if not axis.timestampList.timestamps: + raise ValueError(f"DataBucket for PV '{bucket.pvName}' has a time axis describing no timestamps") + return to_epoch_nanos(axis.timestampList.timestamps[0]) + + +def buckets_by_pv(buckets: Iterable[common_pb2.DataBucket]) -> dict[str, list[common_pb2.DataBucket]]: + """ + Groups buckets by PV name, each PV's buckets ordered by first timestamp. + + PVs keep the order they are first seen in. The sort is stable, so buckets with the same first timestamp keep + their input order. Today's server already returns buckets sorted by (pvName, first time), but does not promise + to, so the grouping is done here; it also reassembles a PV whose buckets were split across pages or stream + messages, once those are combined. + + Overlapping or duplicate timestamps across a PV's buckets are kept, not deduplicated: choosing between two + ingests is not the client's call. + + :param buckets: The buckets, in any order (e.g. the data_buckets of every page). + :return: A dict of PV name to that PV's buckets. + :raises ValueError: if a bucket has no usable time axis. + """ + groups: dict[str, list[common_pb2.DataBucket]] = {} + for bucket in buckets: + groups.setdefault(bucket.pvName, []).append(bucket) + return {pv_name: sorted(group, key=_first_nanos) for pv_name, group in groups.items()} + + +def _structure(column: Any) -> tuple: + """ + The fields of a column, other than its values, without which its values cannot be interpreted: its kind, plus + an enum's enumId, an array's dims, an image's descriptor, or a struct's schemaId. + """ + return ( + type(column).__name__, + column.enumId if isinstance(column, common_pb2.EnumColumn) else None, + tuple(dims) if (dims := dfc.column_dimensions(column)) is not None else None, + tuple(sorted(descriptor.items())) if (descriptor := dfc.image_descriptor_dict(column)) is not None else None, + dfc.column_schema_id(column), + ) + + +def _check_consistent(pv_name: str, buckets: list[common_pb2.DataBucket]) -> None: + """ + Rejects a PV whose buckets differ in column kind or structure, which cannot be concatenated into one column. + + Concatenating anyway would silently widen dtypes (double then int32) or mix payloads that mean different things + (two enumerations, two array shapes). + + :raises ValueError: naming the PV and the first timestamps of the two buckets that differ. + """ + first = buckets[0] + expected = _structure(bucket_column(first)) + for bucket in buckets[1:]: + actual = _structure(bucket_column(bucket)) + if actual != expected: + raise ValueError( + f"PV '{pv_name}' has buckets with different column kinds or structure (the bucket starting at " + f"{_first_nanos(first)} ns has {expected}, the one starting at {_first_nanos(bucket)} ns has " + f"{actual}), so they cannot form one column; read them separately with bucket_to_data_frame() or " + f"bucket_values()" + ) + + +def _bucket_descriptor(bucket: common_pb2.DataBucket, epoch_nanos: list[int], exclude_column_metadata: bool) -> dict: + """The per-bucket entry of df.attrs["buckets"].""" + descriptor: dict[str, Any] = { + "first_nanos": epoch_nanos[0], + "last_nanos": epoch_nanos[-1], + "sample_count": len(epoch_nanos), + "provider_id": bucket.providerId, + "provider_name": bucket.providerName, + } + if not exclude_column_metadata: + descriptor["column_metadata"] = dfc.column_metadata_dict(bucket_column(bucket)) + return descriptor + + +def _pv_dataframe(pd: Any, pv_name: str, buckets: list[common_pb2.DataBucket], exclude_column_metadata: bool) -> Any: + """Builds one PV's pandas DataFrame from its (already trimmed, ordered, consistent) buckets.""" + parts = [] + descriptors = [] + for bucket in buckets: + frame = bucket_to_data_frame(bucket) + part = dfc.data_frame_to_pandas(frame, exclude_column_metadata=True) + # The column is named for the ingested column, which ingestion makes the PV name; name it for the PV + # regardless, so every part concatenates into the one column. + part.columns = [pv_name] + parts.append(part) + descriptors.append(_bucket_descriptor(bucket, dfc.data_frame_timestamps(frame), exclude_column_metadata)) + + df = pd.concat(parts, axis=0) if len(parts) > 1 else parts[0] + + # Rebuilt from scratch rather than inherited: how pd.concat propagates attrs has varied across pandas versions. + column = bucket_column(buckets[0]) + attrs: dict[str, Any] = {"buckets": descriptors} + if isinstance(column, common_pb2.EnumColumn): + # Structural, not metadata: carried even under exclude_column_metadata, as data_frame_to_pandas() does. + attrs["enum_ids"] = {pv_name: column.enumId} + dims = dfc.column_dimensions(column) + if dims is not None: + attrs["dimensions"] = {pv_name: dims} + descriptor = dfc.image_descriptor_dict(column) + if descriptor is not None: + attrs["image_descriptors"] = {pv_name: descriptor} + schema_id = dfc.column_schema_id(column) + if schema_id is not None: + attrs["schema_ids"] = {pv_name: schema_id} + if not exclude_column_metadata: + # Each ingest carries its own metadata, so a PV's buckets can legitimately differ. The per-PV summary is set + # only when they agree; the per-bucket copies in attrs["buckets"] are always there. + first_metadata = descriptors[0]["column_metadata"] + if all(entry["column_metadata"] == first_metadata for entry in descriptors[1:]): + attrs["column_metadata"] = {pv_name: first_metadata} + df.attrs = attrs + return df + + +def buckets_to_dataframes( + buckets: Iterable[common_pb2.DataBucket], + *, + time_range: tuple[TimestampInput, TimestampInput] | None = None, + exclude_column_metadata: bool = False, +) -> dict[str, Any]: + """ + Assembles buckets into one pandas DataFrame per PV. Requires the optional [analysis] extra. + + Each frame has a UTC DatetimeIndex built from int64 epoch nanoseconds and one column, named for the PV, with the + dtype its column type implies (float32 for a FloatColumn, int32 for an Int32Column or EnumColumn, ...). Array, + image, and struct samples are object cells, as in data_frame_conversions.data_frame_to_pandas(). A PV's buckets + are concatenated in first-timestamp order (see buckets_by_pv()); overlapping buckets are kept, so the index can + be non-monotonic or carry repeated instants. + + df.attrs carries: + - "buckets": one dict per bucket in the frame, with 'first_nanos', 'last_nanos', 'sample_count', + 'provider_id', 'provider_name', and (unless exclude_column_metadata) 'column_metadata'. + - "column_metadata": {pv: dict}, only when every bucket's metadata is identical and not excluded. + - "enum_ids", "dimensions", "image_descriptors", "schema_ids": {pv: value} for the column kinds that carry + them. These are structural -- the values cannot be interpreted without them -- so they are present even + under exclude_column_metadata. + + :param buckets: The buckets, in any order (e.g. the data_buckets of every page or stream message). + :param time_range: Optional (begin, end) to trim each bucket to [begin, end) exactly (see trim_bucket()); None + (the default) keeps every bucket whole, as the server returned it. A PV with no sample left in the range is + absent from the result. + :param exclude_column_metadata: When True, leave column metadata out of attrs. + :return: A dict of PV name to pandas.DataFrame, in first-seen PV order. + :raises ImportError: if the [analysis] extra is not installed. + :raises ValueError: if time_range's begin is not before its end; if a PV's buckets differ in column kind or + structure; or if a bucket is malformed, holds a SerializedDataColumn, or (when trimming) has a decreasing + time axis. + """ + pd = _require_pandas() + + range_nanos = _range_nanos(*time_range) if time_range is not None else None + + frames: dict[str, Any] = {} + for pv_name, group in buckets_by_pv(buckets).items(): + for bucket in group: + _refuse_serialized(bucket, bucket_column(bucket), "converted to a pandas DataFrame") + if range_nanos is not None: + group = [trimmed for bucket in group if (trimmed := _trim_nanos(bucket, *range_nanos)) is not None] + if not group: + continue + _check_consistent(pv_name, group) + frames[pv_name] = _pv_dataframe(pd, pv_name, group, exclude_column_metadata) + return frames + + +def query_buckets_to_dataframes( + query_client: Any, params: Any, *, trim: bool = False, max_buckets: int | None = None +) -> dict[str, Any]: + """ + Runs a unary bucket query, pages through the whole result with iter_query_buckets(), and assembles one pandas + DataFrame per PV (see buckets_to_dataframes()). Requires the optional [analysis] extra. + + There is deliberately no streaming counterpart: a PV can span stream messages, so a frame per message would be + a PV fragment. Stream callers collect the buckets of every message and call buckets_to_dataframes() once. + + :param query_client: A QueryClient. + :param params: The QueryParams describing the query (must not carry a sample_status_filter). + :param trim: When True, trim every bucket to the query's own [begin, end); by default buckets stay whole. + :param max_buckets: Optional cap on the total number of buckets across all pages; raise once it would be + exceeded (an unbounded "give me everything" is an out-of-memory foot-gun on large ranges). + :return: A dict of PV name to pandas.DataFrame. + :raises ValueError: if the result exceeds max_buckets, or the buckets cannot be assembled. + :raises RuntimeError: if any page returns an error (propagated from iter_query_buckets()). + """ + _require_pandas() + + buckets: list[common_pb2.DataBucket] = [] + for page in query_client.iter_query_buckets(params): + buckets.extend(page.data_buckets) + if max_buckets is not None and len(buckets) > max_buckets: + raise ValueError( + f"query result exceeds max_buckets={max_buckets} (at least {len(buckets)} buckets); narrow the " + f"range or raise max_buckets" + ) + + time_range = (params.begin_timestamp, params.end_timestamp) if trim else None + return buckets_to_dataframes(buckets, time_range=time_range, exclude_column_metadata=params.exclude_column_metadata) diff --git a/src/dp_python_lib/client/data_frame_conversions.py b/src/dp_python_lib/client/data_frame_conversions.py index a268425..472d800 100644 --- a/src/dp_python_lib/client/data_frame_conversions.py +++ b/src/dp_python_lib/client/data_frame_conversions.py @@ -1,8 +1,8 @@ """ Pythonic conversions for common.DataFrame (issue #6, Phase 2 and 3). -Reads a DataFrame -- an annotation's calculations, and in future a bucket query's typed columns -- back into plain -Python, and, behind the optional [analysis] extra, into pandas. +Reads a DataFrame -- an annotation's calculations, or a bucket query's bucket viewed as a one-column frame (see +bucket_conversions, #16) -- back into plain Python, and, behind the optional [analysis] extra, into pandas. The pure-Python half has no third-party dependencies. pandas is imported lazily inside each entry point that needs it, so importing this module never requires the extra. @@ -11,8 +11,8 @@ - Timestamps are computed in INTEGER NANOSECONDS via expand_data_timestamps(), never in float seconds. A float64 carries 53 bits of mantissa and present-day epoch nanoseconds need ~61, so a float round-trip would silently move every timestamp. The pandas index is built from those int64 nanoseconds directly for the same reason. - - column_values() is written as a standalone per-column converter because the bucket query (#16) reuses these - same 14 typed column messages; it takes one column and needs to know nothing about the frame. + - column_values() is written as a standalone per-column converter because the bucket query (#16) carries these + same typed column messages one per bucket; it takes one column and needs to know nothing about the frame. - Array columns are reshaped into one list per sample using their declared dims, rather than returned flat: a flat list would silently lose the sample boundaries. - Dense typed columns cannot express a gap, so the pandas->DataFrame direction rejects NaN/None fail-loud with a @@ -235,8 +235,8 @@ def column_values(column: Any) -> list: """ Extracts one Python value per sample from any supported column message. - Written as a standalone converter because the bucket query (#16) carries the same 14 typed column messages and - can reuse this without going through a DataFrame. + Written as a standalone converter because a bucket query result (#16) carries the same typed column messages, + one per bucket; bucket_conversions.bucket_values() reuses it without going through a DataFrame. Mapping: typed scalar columns yield their native values; an EnumColumn yields its integer codes (the enumeration naming them is `enumId` on the column); an array column yields one list per sample; an ImageColumn yields one diff --git a/src/dp_python_lib/client/query_client.py b/src/dp_python_lib/client/query_client.py index d342c95..ef85d74 100644 --- a/src/dp_python_lib/client/query_client.py +++ b/src/dp_python_lib/client/query_client.py @@ -1,6 +1,6 @@ import logging -from collections.abc import Iterator -from typing import Any +from collections.abc import Callable, Iterable, Iterator +from typing import Any, TypeVar import grpc @@ -9,6 +9,9 @@ from dp_python_lib.client.time_conversions import TimestampInput, to_timestamp from dp_python_lib.grpc import common_pb2, query_pb2, query_pb2_grpc +# The result type a server-streaming sender yields (QuerySamplesApiResult or QueryBucketsApiResult). +_StreamResultT = TypeVar("_StreamResultT", bound=ApiResultBase) + class PvQuery: """ @@ -353,8 +356,9 @@ def exclude( class QueryParams: """ Encapsulates client parameters for a v2 time-series query. This is the kind-neutral representation of the shared - QuerySpec: it is used to build both the sample-oriented querySamples()/querySamplesStream() requests (in scope for - this release) and, in a future release, the bucket-oriented queryBuckets() requests. + QuerySpec: it builds both the sample-oriented querySamples()/querySamplesStream() requests and the bucket-oriented + queryBuckets()/queryBucketsStream() requests (see QueryClient), so a caller switching kinds changes only the + method name. A query selects data over a half-open time range [begin_time, end_time) for a set of PVs. The PV set is chosen by exactly one pv_selector form (see PvQuery: name-list, pattern, or metadata), optionally restricted by a list @@ -382,13 +386,18 @@ def __init__( form set -- the PvQuery helpers each produce exactly one form. For every PV, use PvQuery.pattern(".*"). :param config_criteria: List of AND-combined configuration criteria (see ConfigQuery) restricting results to intervals when matching configurations were active. Optional. - :param limit: Maximum number of rows to return per page (optional). This is a per-page size, NOT a total - cap. 0 is meaningful and means "let the server pick a default"; a negative value raises. + :param limit: Maximum number of results per page (optional). This is a per-page size, NOT a total cap, + and what it counts depends on the query kind: ROWS for querySamples(), BUCKETS for queryBuckets(). A + bucket page can also end early, with a page token, when it reaches the server's outbound byte budget, so + it may hold fewer than limit buckets while more remain; and the server silently clamps a bucket limit + to its configured maximum (100,000 by default). 0 is meaningful and means "let the server pick a + default"; a negative value raises. :param exclude_column_metadata: If True, omit per-column ColumnMetadata from the results. Defaults to False (metadata included). :param sample_status_filter: Optional SampleStatusSelector (see SampleStatusFilter.include()/exclude()) restricting results to, or away from, samples carrying matching sample statuses. Supported by - querySamples()/querySamplesStream() only -- see the note in _build_query_spec(). + querySamples()/querySamplesStream() only: the bucket methods refuse a QueryParams carrying one with a + ValueError, since the server rejects it on a bucket query. :raises ValueError: if both begin_time and end_time are not supplied, if pv_selector is None or sets no selector form, if begin_time is not strictly before end_time, if limit is negative, or if sample_status_filter is present but carries an empty domain or MODE_UNSPECIFIED. @@ -520,12 +529,88 @@ def to_numpy(self) -> Any: return query_conversions.column_table_to_numpy(self.column_table) +class QueryBucketsApiResult(ApiResultBase): + """ + Wraps a single page (unary) or a single streamed message (streaming) of a queryBuckets()/queryBucketsStream() + response, with a status object including an error flag and message. + + The raw protobuf DataBuckets are available via .data_buckets. Each is one stored unit of the archive, returned + WHOLE: the server selects every bucket overlapping [begin, end) and never trims it, so a boundary bucket carries + samples outside the query range. See bucket_conversions for reading, trimming, and assembling them. + + There is deliberately no .to_dataframe(): a page of buckets is not one table. .to_dataframes() returns one + pandas DataFrame per PV instead, and requires the optional [analysis] extra. + """ + + def __init__( + self, + is_error: bool, + message: str, + response: query_pb2.QueryBucketsResponse | None = None, + ) -> None: + """ + :param is_error: Boolean flag indicating if an error occurred in the API call. + :param message: Error message describing the error condition. + :param response: The QueryBucketsResponse for this page/message, or None. + """ + super().__init__(is_error, message) + self.response = response + + @property + def data_buckets(self) -> list[common_pb2.DataBucket]: + """The DataBuckets in this page/message, in server order; an empty list on error or for an empty result.""" + if self.response is not None and self.response.HasField("bucketQueryResult"): + return list(self.response.bucketQueryResult.dataBuckets) + return [] + + @property + def next_page_token(self) -> str: + """ + Token for retrieving the next page (unary queryBuckets() only), or empty string if there are no more pages. + Always empty for streamed messages (queryBucketsStream() is fire-and-consume). + """ + if self.response is not None and self.response.HasField("bucketQueryResult"): + return self.response.bucketQueryResult.nextPageToken + return "" + + def to_dataframes( + self, + time_range: tuple[TimestampInput, TimestampInput] | None = None, + exclude_column_metadata: bool = False, + ) -> dict[str, Any]: + """ + Converts this page's buckets into one pandas DataFrame per PV. Requires the optional [analysis] extra. + Delegates to bucket_conversions.buckets_to_dataframes() (imported lazily so the core client carries no + pandas dependency); see it for the frame shape and attrs. + + A PV's buckets can span pages, so converting page by page yields that PV's buckets in pieces. Collect the + buckets of every page and convert once, or use bucket_conversions.query_buckets_to_dataframes(). + + :param time_range: Optional (begin, end) to trim each bucket to [begin, end) exactly; None (the default) + leaves the buckets whole, as the server returned them. + :param exclude_column_metadata: If True, do not attach per-bucket ColumnMetadata to the DataFrames. + :return: A dict of PV name to pandas.DataFrame. + """ + from dp_python_lib.client import bucket_conversions + + return bucket_conversions.buckets_to_dataframes( + self.data_buckets, time_range=time_range, exclude_column_metadata=exclude_column_metadata + ) + + class QueryClient(ServiceApiClientBase): """ - User-facing client for the sample-oriented v2 query methods of the MLDP Query Service: querySamples() (unary, one - resumable page) and querySamplesStream() (server-streaming, fire-and-consume). Provides low-level wrappers around - the raw protobuf ColumnTable plus transparent-paging iterators. Higher-level pandas/NumPy/Excel conversions live in - query_conversions and are reached via QuerySamplesApiResult.to_dataframe()/.to_numpy(). + User-facing client for the v2 query methods of the MLDP Query Service, in two kinds: + + - sample-oriented: querySamples() (unary, one resumable page) and querySamplesStream() (server-streaming, + fire-and-consume), returning an aligned, trimmed, scalar-only ColumnTable. Pandas/NumPy/Excel conversions + live in query_conversions and are reached via QuerySamplesApiResult.to_dataframe()/.to_numpy(). + - bucket-oriented: queryBuckets() and queryBucketsStream(), returning the archive's stored DataBuckets whole, + each with its original typed column (any kind, arrays and images included) and time axis. Conversions live + in bucket_conversions and are reached via QueryBucketsApiResult.to_dataframes(). + + Each kind has the same three methods with the same contracts: query_*() for one page, iter_query_*() to page + transparently, and iter_query_*_stream() for the server stream. Queries are described by a kind-neutral QueryParams built from the PvQuery (PV) and ConfigQuery (CFG) helpers. """ @@ -545,14 +630,12 @@ def __init__(self, channel: grpc.Channel) -> None: def _build_query_spec(self, request_params: QueryParams) -> query_pb2.QuerySpec: """ Builds the shared QuerySpec (time range + PV selector + configuration selector + optional sample status - selector) from the supplied QueryParams. Factored out so both _build_query_samples_request() and a future - _build_query_buckets_request() reuse it. + selector) from the supplied QueryParams. Shared by _build_query_samples_request() and + _build_query_buckets_request(). - NOTE for the future bucket-oriented client (issue #16): sampleStatusSelector is supported by - querySamples()/querySamplesStream() ONLY -- the server rejects it on a bucket query, because a bucket is a - stored unit that cannot be partially filtered without rewriting it. A bucket request builder reusing this - method must therefore refuse a QueryParams carrying sample_status_filter rather than copying it through, - which would silently produce a request the server rejects. + This copies sample_status_filter through whenever it is set. That is right for samples only: the server + rejects a sampleStatusSelector on a bucket query, so _build_query_buckets_request() refuses such params + before calling here. Any other caller added later must make the same choice. :param request_params: User parameters for the query. :return: A QuerySpec for the specified params. @@ -599,6 +682,102 @@ def _build_query_samples_request( return request + def _build_query_buckets_request( + self, request_params: QueryParams, page_token: str | None = None + ) -> query_pb2.QueryBucketsRequest: + """ + Builds a QueryBucketsRequest from the supplied QueryParams and optional page token. Used by both the unary and + streaming RPCs (they share the request type); the streaming path must never supply a page_token. + + :param request_params: User parameters for the query. + :param page_token: Token for retrieving a subsequent page (unary paging only). + :return: A QueryBucketsRequest for the specified params. + :raises ValueError: if request_params carries a sample_status_filter, which the server rejects on a bucket + query. + """ + # Refused rather than dropped: dropping it would hand back unfiltered data to a caller who believes it was + # filtered. A bucket is a stored unit that cannot be partially filtered without rewriting it, and the + # server rejects the selector's mere presence (dp-service QueryV2Resolver). + if request_params.sample_status_filter is not None: + raise ValueError( + "bucket-oriented queries do not support per-sample status filtering; remove sample_status_filter " + "or use query_samples()/iter_query_samples()/iter_query_samples_stream()" + ) + + self.logger.debug("Building QueryBucketsRequest") + request = query_pb2.QueryBucketsRequest() + request.querySpec.CopyFrom(self._build_query_spec(request_params)) + + if request_params.limit is not None: + request.executionOptions.limit = request_params.limit + if page_token: + request.executionOptions.pageToken = page_token + + # useSerializedColumns is inert on buckets (dp-service passes it in and never reads it): a column comes back + # serialized only if it was stored that way. Forced False so the request never asks for something the server + # does not do. excludeColumnMetadata, unlike on the samples path, is honored here. + request.resultRepresentation.useSerializedColumns = False + request.resultRepresentation.excludeColumnMetadata = request_params.exclude_column_metadata + + return request + + # ------------------------------------------------------------------ + # server-streaming (shared by both query kinds) + # ------------------------------------------------------------------ + + def _iter_stream( + self, + stub_call: Callable[[Any], Iterable[Any]], + request: Any, + result_cls: Callable[..., _StreamResultT], + success_field: str, + op_name: str, + ) -> Iterator[_StreamResultT]: + """ + Invokes a server-streaming query method, yielding one result per streamed message. + + Errors are yielded, not raised: a business error on a message, an unrecognized response, or a gRPC/unexpected + error while iterating the stream each yield an is_error result, and a transport error terminates the + generator after that final error item. This is the internal contract -- the public iter_*_stream() wrappers + consume these and convert the first error result into a RuntimeError, so callers of the public methods never + see an error result yielded. + + Deliberately not routed through _dispatch(), which returns a single result rather than yielding many. + + :param stub_call: The bound stub method (e.g. self._stub.querySamplesStream). + :param request: The request for the call (must carry no page token). + :param result_cls: The result class; called as result_cls(is_error=, message=, response=). + :param success_field: The success oneof field of each response (e.g. "sampleQueryResult"). + :param op_name: The camelCase method name, for log and error messages. + :return: An iterator over the streamed result messages, possibly ending in an error result. + """ + display_name = op_name[0].upper() + op_name[1:] + self.logger.info("Calling %s API", op_name) + + try: + self.logger.debug("Invoking stub.%s with request", op_name) + stream = stub_call(request) + for response in stream: + if response.HasField("exceptionalResult"): + error_msg = response.exceptionalResult.message + self.logger.warning("%s returned business error: %s", display_name, error_msg) + yield result_cls(is_error=True, message=error_msg) + elif response.HasField(success_field): + yield result_cls(is_error=False, message="", response=response) + else: + error_msg = f"Unexpected response format: neither exceptionalResult nor {success_field} found" + self.logger.error(error_msg) + yield result_cls(is_error=True, message=error_msg) + + except grpc.RpcError as e: + error_msg = f"gRPC error: {e.details()}" + self.logger.error("gRPC error during %s: %s", op_name, e.details()) + yield result_cls(is_error=True, message=error_msg) + except Exception as e: + error_msg = f"Unexpected error: {e!s}" + self.logger.exception("Unexpected error during %s: %s", op_name, str(e)) + yield result_cls(is_error=True, message=error_msg) + # ------------------------------------------------------------------ # querySamples (unary) # ------------------------------------------------------------------ @@ -668,42 +847,18 @@ def iter_query_samples(self, request_params: QueryParams) -> Iterator[QuerySampl def _send_query_samples_stream(self, request: query_pb2.QuerySamplesRequest) -> Iterator[QuerySamplesApiResult]: """ Invokes the querySamplesStream() server-streaming API method with the supplied request, yielding one - QuerySamplesApiResult per streamed message. - - Errors are yielded, not raised: a business error on a message, an unrecognized response, or a gRPC/unexpected - error while iterating the stream each yield an is_error result, and a transport error terminates the - generator after that final error item. This is the internal contract -- the public - iter_query_samples_stream() wrapper consumes these and converts the first error result into a RuntimeError, - so callers of the public method never see an error result yielded. + QuerySamplesApiResult per streamed message. Errors are yielded, not raised; see _iter_stream(). :param request: QuerySamplesRequest with parameters for the call (must carry no page token). :return: An iterator over the streamed result messages, possibly ending in an error result. """ - self.logger.info("Calling querySamplesStream API") - - try: - self.logger.debug("Invoking stub.querySamplesStream with request") - stream = self._stub.querySamplesStream(request) - for response in stream: - if response.HasField("exceptionalResult"): - error_msg = response.exceptionalResult.message - self.logger.warning("QuerySamplesStream returned business error: %s", error_msg) - yield QuerySamplesApiResult(is_error=True, message=error_msg) - elif response.HasField("sampleQueryResult"): - yield QuerySamplesApiResult(is_error=False, message="", response=response) - else: - error_msg = "Unexpected response format: neither exceptionalResult nor sampleQueryResult found" - self.logger.error(error_msg) - yield QuerySamplesApiResult(is_error=True, message=error_msg) - - except grpc.RpcError as e: - error_msg = f"gRPC error: {e.details()}" - self.logger.error("gRPC error during querySamplesStream: %s", e.details()) - yield QuerySamplesApiResult(is_error=True, message=error_msg) - except Exception as e: - error_msg = f"Unexpected error: {e!s}" - self.logger.exception("Unexpected error during querySamplesStream: %s", str(e)) - yield QuerySamplesApiResult(is_error=True, message=error_msg) + return self._iter_stream( + self._stub.querySamplesStream, + request, + QuerySamplesApiResult, + "sampleQueryResult", + "querySamplesStream", + ) def iter_query_samples_stream(self, request_params: QueryParams) -> Iterator[QuerySamplesApiResult]: """ @@ -728,3 +883,120 @@ def iter_query_samples_stream(self, request_params: QueryParams) -> Iterator[Que if result.result_status.is_error: raise RuntimeError(f"querySamplesStream failed during streaming: {result.result_status.message}") yield result + + # ------------------------------------------------------------------ + # queryBuckets (unary) + # ------------------------------------------------------------------ + + def _send_query_buckets(self, request: query_pb2.QueryBucketsRequest) -> QueryBucketsApiResult: + """ + Invokes the queryBuckets() unary API method with the supplied request. + :param request: QueryBucketsRequest with parameters for the call. + :return: A QueryBucketsApiResult with the method response and status information. + """ + return self._dispatch( + self._stub.queryBuckets, + request, + QueryBucketsApiResult, + "bucketQueryResult", + "queryBuckets", + success_log=lambda response: self.logger.info( + "QueryBuckets returned %d buckets", len(response.bucketQueryResult.dataBuckets) + ), + ) + + def query_buckets(self, request_params: QueryParams, page_token: str | None = None) -> QueryBucketsApiResult: + """ + User-facing method for invoking the unary queryBuckets() API method. Returns a single page of whole buckets; + use iter_query_buckets() to page through all results transparently. + + params.limit counts BUCKETS here, not rows, and a page may also end early (with a page token) at the server's + outbound byte budget. An empty result is a success with no buckets. + + :param request_params: User parameters for the query (see QueryParams / PvQuery / ConfigQuery). Must not + carry a sample_status_filter. + :param page_token: Token for retrieving a subsequent page (optional). + :return: A QueryBucketsApiResult with a single page of results and status information. + :raises ValueError: if request_params carries a sample_status_filter. + """ + self.logger.info("Starting queryBuckets operation") + + request = self._build_query_buckets_request(request_params, page_token=page_token) + result = self._send_query_buckets(request) + + if result.result_status.is_error: + self.logger.error("QueryBuckets operation failed: %s", result.result_status.message) + else: + self.logger.info("QueryBuckets operation completed successfully") + + return result + + def iter_query_buckets(self, request_params: QueryParams) -> Iterator[QueryBucketsApiResult]: + """ + Convenience generator that transparently pages through all unary queryBuckets() results, following the + nextPageToken until the results are exhausted. Yields one QueryBucketsApiResult per page. + + Raises RuntimeError if any page returns an error, so callers can distinguish failure from an empty result set. + + :param request_params: User parameters for the query (see QueryParams / PvQuery / ConfigQuery). Must not + carry a sample_status_filter. + :return: An iterator over the result pages. + :raises ValueError: if request_params carries a sample_status_filter (on the first next()). + """ + page_token: str | None = None + while True: + result = self.query_buckets(request_params, page_token=page_token) + if result.result_status.is_error: + raise RuntimeError(f"queryBuckets failed during paging: {result.result_status.message}") + + yield result + + page_token = result.next_page_token + if not page_token: + break + + # ------------------------------------------------------------------ + # queryBucketsStream (server-streaming) + # ------------------------------------------------------------------ + + def _send_query_buckets_stream(self, request: query_pb2.QueryBucketsRequest) -> Iterator[QueryBucketsApiResult]: + """ + Invokes the queryBucketsStream() server-streaming API method with the supplied request, yielding one + QueryBucketsApiResult per streamed message. Errors are yielded, not raised; see _iter_stream(). + + :param request: QueryBucketsRequest with parameters for the call (must carry no page token). + :return: An iterator over the streamed result messages, possibly ending in an error result. + """ + return self._iter_stream( + self._stub.queryBucketsStream, + request, + QueryBucketsApiResult, + "bucketQueryResult", + "queryBucketsStream", + ) + + def iter_query_buckets_stream(self, request_params: QueryParams) -> Iterator[QueryBucketsApiResult]: + """ + User-facing lazy generator for the server-streaming queryBucketsStream() API method. Yields one + QueryBucketsApiResult per streamed message, symmetric with iter_query_buckets() but fire-and-consume: the + server reads the whole result and pushes it in messages cut by params.limit (in buckets) and by its byte + budget, with no page tokens. An empty result is one message with no buckets. + + A PV's buckets can span messages, so per-message conversion yields PV fragments; collect the buckets and + call bucket_conversions.buckets_to_dataframes() once to get whole PVs. + + Raises RuntimeError on a mid-stream error, matching the iter_* page-error behavior, so callers can distinguish + failure from an empty stream. + + :param request_params: User parameters for the query (see QueryParams / PvQuery / ConfigQuery). Must not + carry a sample_status_filter. + :return: A lazy iterator over the streamed result messages. + :raises ValueError: if request_params carries a sample_status_filter (on the first next()). + """ + self.logger.info("Starting queryBucketsStream operation") + + request = self._build_query_buckets_request(request_params, page_token=None) + for result in self._send_query_buckets_stream(request): + if result.result_status.is_error: + raise RuntimeError(f"queryBucketsStream failed during streaming: {result.result_status.message}") + yield result diff --git a/tests/unit/test_bucket_conversions.py b/tests/unit/test_bucket_conversions.py new file mode 100644 index 0000000..15c81aa --- /dev/null +++ b/tests/unit/test_bucket_conversions.py @@ -0,0 +1,542 @@ +import os +import sys +import unittest +from datetime import datetime, timezone +from unittest.mock import Mock + +# Add src directory to path for imports +sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) + +from dp_python_lib.client import bucket_conversions as bc +from dp_python_lib.client import data_frame as dfb +from dp_python_lib.client.query_client import PvQuery, QueryBucketsApiResult, QueryParams +from dp_python_lib.client.time_conversions import from_epoch_nanos, to_epoch_nanos +from dp_python_lib.grpc import common_pb2, query_pb2 + +# The pandas half depends on the optional [analysis] extra; skip those tests cleanly when it is absent. The +# pure-Python half needs no optional deps and always runs. +try: + import pandas # noqa: F401 -- availability probe for the [analysis] extra + + _HAVE_ANALYSIS = True +except ImportError: + _HAVE_ANALYSIS = False + +# A present-day epoch: its nanoseconds need ~61 bits, so any float round trip would move them. +T0 = datetime(2026, 9, 30, 18, 0, 0, 123456, tzinfo=timezone.utc) +T0_NANOS = to_epoch_nanos(dfb.sampling_clock(T0, 1, 1).samplingClock.startTime) +PERIOD = 1_000_000 + + +def _clock(count=5, start_nanos=T0_NANOS, period=PERIOD): + return dfb.sampling_clock(from_epoch_nanos(start_nanos), period, count) + + +def _raw_list(epoch_nanos): + """A TimestampList built without timestamp_list(), which refuses the repeated/decreasing axes some tests need.""" + axis = common_pb2.DataTimestamps() + for nanos in epoch_nanos: + axis.timestampList.timestamps.add().CopyFrom(from_epoch_nanos(nanos)) + return axis + + +def _bucket(pv, axis, column, provider_id="prov-1", provider_name="daq"): + bucket = common_pb2.DataBucket(pvName=pv, providerId=provider_id, providerName=provider_name) + bucket.dataTimestamps.CopyFrom(axis) + arm = { + common_pb2.DataColumn: "dataColumn", + common_pb2.SerializedDataColumn: "serializedDataColumn", + common_pb2.DoubleColumn: "doubleColumn", + common_pb2.FloatColumn: "floatColumn", + common_pb2.Int64Column: "int64Column", + common_pb2.Int32Column: "int32Column", + common_pb2.BoolColumn: "boolColumn", + common_pb2.StringColumn: "stringColumn", + common_pb2.EnumColumn: "enumColumn", + common_pb2.ImageColumn: "imageColumn", + common_pb2.StructColumn: "structColumn", + common_pb2.DoubleArrayColumn: "doubleArrayColumn", + common_pb2.FloatArrayColumn: "floatArrayColumn", + common_pb2.Int32ArrayColumn: "int32ArrayColumn", + common_pb2.Int64ArrayColumn: "int64ArrayColumn", + common_pb2.BoolArrayColumn: "boolArrayColumn", + }[type(column)] + getattr(bucket.dataValues, arm).CopyFrom(column) + return bucket + + +def _double_bucket(pv="PV:A", count=5, start_nanos=T0_NANOS, values=None, **kwargs): + values = values if values is not None else [float(i) for i in range(count)] + return _bucket(pv, _clock(count, start_nanos), dfb.double_column(pv, values), **kwargs) + + +# Every column arm a DataBucket can carry except the serialized one, with three samples apiece and the values +# bucket_values() should yield. +def _all_arm_cases(pv="PV:X"): + return [ + (dfb.double_column(pv, [1.5, 2.5, 3.5]), [1.5, 2.5, 3.5]), + (dfb.float_column(pv, [0.5, 1.5, 2.5]), [0.5, 1.5, 2.5]), + (dfb.int64_column(pv, [2**40, -1, 0]), [2**40, -1, 0]), + (dfb.int32_column(pv, [1, 2, 3]), [1, 2, 3]), + (dfb.bool_column(pv, [True, False, True]), [True, False, True]), + (dfb.string_column(pv, ["a", "b", "c"]), ["a", "b", "c"]), + (dfb.enum_column(pv, [0, 2, 1], "mode:v1"), [0, 2, 1]), + (dfb.image_column(pv, [b"i0", b"i1", b"i2"], 2, 1, 1, "raw-u8"), [b"i0", b"i1", b"i2"]), + (dfb.struct_column(pv, [b"s0", b"s1", b"s2"], "schema:v1"), [b"s0", b"s1", b"s2"]), + (dfb.double_array_column(pv, [[1.0, 2.0], [3.0, 4.0], [5.0, 6.0]]), [[1.0, 2.0], [3.0, 4.0], [5.0, 6.0]]), + (dfb.float_array_column(pv, [[0.5], [1.5], [2.5]]), [[0.5], [1.5], [2.5]]), + (dfb.int32_array_column(pv, [[1, 2], [3, 4], [5, 6]]), [[1, 2], [3, 4], [5, 6]]), + (dfb.int64_array_column(pv, [[1], [2], [3]]), [[1], [2], [3]]), + (dfb.bool_array_column(pv, [[True], [False], [True]]), [[True], [False], [True]]), + (dfb.data_column(pv, [1.5, None, "x"]), [1.5, None, "x"]), + ] + + +# ---------------------------------------------------------------------- +# one bucket, pure Python +# ---------------------------------------------------------------------- + + +class TestBucketColumn(unittest.TestCase): + def test_returns_the_set_arm(self): + bucket = _double_bucket() + self.assertIsInstance(bc.bucket_column(bucket), common_pb2.DoubleColumn) + + def test_serialized_column_is_reachable(self): + column = dfb.serialized_column("PV:S", b"\x01\x02", "my-codec") + bucket = _bucket("PV:S", _clock(2), column) + reached = bc.bucket_column(bucket) + self.assertEqual(reached.payload, b"\x01\x02") + self.assertEqual(reached.encoding, "my-codec") + + def test_unset_data_values_raises(self): + bucket = common_pb2.DataBucket(pvName="PV:E") + bucket.dataTimestamps.CopyFrom(_clock(1)) + with self.assertRaises(ValueError) as ctx: + bc.bucket_column(bucket) + self.assertIn("PV:E", str(ctx.exception)) + + +class TestBucketToDataFrame(unittest.TestCase): + def test_every_arm_lands_in_its_field_and_reads_back(self): + from dp_python_lib.client import data_frame_conversions as dfc + + for column, expected in _all_arm_cases(): + with self.subTest(kind=type(column).__name__): + bucket = _bucket("PV:X", _clock(3), column) + frame = bc.bucket_to_data_frame(bucket) + self.assertEqual(dfc.data_frame_columns(frame), {"PV:X": expected}) + self.assertEqual(frame.dataTimestamps, bucket.dataTimestamps) + + def test_serialized_lands_in_serialized_columns(self): + from dp_python_lib.client import data_frame_conversions as dfc + + bucket = _bucket("PV:S", _clock(2), dfb.serialized_column("PV:S", b"xy", "codec")) + frame = bc.bucket_to_data_frame(bucket) + self.assertEqual(len(frame.serializedDataColumns), 1) + self.assertEqual(dfc.data_frame_columns(frame), {}) # skipped by design + + def test_unset_axis_raises(self): + bucket = common_pb2.DataBucket(pvName="PV:N") + bucket.dataValues.doubleColumn.CopyFrom(dfb.double_column("PV:N", [1.0])) + with self.assertRaises(ValueError) as ctx: + bc.bucket_to_data_frame(bucket) + self.assertIn("PV:N", str(ctx.exception)) + self.assertIn("dataTimestamps", str(ctx.exception)) + + +class TestBucketTimestampsAndValues(unittest.TestCase): + def test_sampling_clock_is_nanosecond_exact(self): + bucket = _bucket("PV:A", _clock(3, period=7), dfb.double_column("PV:A", [1.0, 2.0, 3.0])) + self.assertEqual(bc.bucket_timestamps(bucket), [T0_NANOS, T0_NANOS + 7, T0_NANOS + 14]) + + def test_timestamp_list_is_nanosecond_exact(self): + nanos = [T0_NANOS + 1, T0_NANOS + 999, T0_NANOS + 1_000_000_001] + bucket = _bucket("PV:A", _raw_list(nanos), dfb.double_column("PV:A", [1.0, 2.0, 3.0])) + self.assertEqual(bc.bucket_timestamps(bucket), nanos) + + def test_empty_list_axis_raises(self): + bucket = _bucket("PV:A", _raw_list([]), dfb.double_column("PV:A", [1.0])) + bucket.dataTimestamps.timestampList.SetInParent() + with self.assertRaises(ValueError): + bc.bucket_timestamps(bucket) + + def test_every_arm_values(self): + for column, expected in _all_arm_cases(): + with self.subTest(kind=type(column).__name__): + self.assertEqual(bc.bucket_values(_bucket("PV:X", _clock(3), column)), expected) + + def test_legacy_gap_is_none(self): + bucket = _bucket("PV:G", _clock(3), dfb.data_column("PV:G", [1, None, 3])) + self.assertEqual(bc.bucket_values(bucket), [1, None, 3]) + + def test_serialized_values_refused(self): + bucket = _bucket("PV:S", _clock(2), dfb.serialized_column("PV:S", b"xy", "my-codec")) + with self.assertRaises(ValueError) as ctx: + bc.bucket_values(bucket) + self.assertIn("PV:S", str(ctx.exception)) + self.assertIn("my-codec", str(ctx.exception)) + self.assertIn("bucket_column()", str(ctx.exception)) + + def test_misaligned_bucket_raises(self): + bucket = _bucket("PV:M", _clock(3), dfb.double_column("PV:M", [1.0, 2.0])) + with self.assertRaises(ValueError) as ctx: + bc.bucket_values(bucket) + self.assertIn("PV:M", str(ctx.exception)) + + def test_repeated_timestamps_convert_without_error(self): + # Ingestion allows a non-decreasing list; the client's write-side strictly-increasing rule must not apply. + nanos = [T0_NANOS, T0_NANOS, T0_NANOS + 5] + bucket = _bucket("PV:D", _raw_list(nanos), dfb.double_column("PV:D", [1.0, 2.0, 3.0])) + self.assertEqual(bc.bucket_timestamps(bucket), nanos) + self.assertEqual(bc.bucket_values(bucket), [1.0, 2.0, 3.0]) + + +# ---------------------------------------------------------------------- +# trimming +# ---------------------------------------------------------------------- + + +class TestTrimBucket(unittest.TestCase): + def _ns(self, offset): + return from_epoch_nanos(T0_NANOS + offset) + + def test_begin_kept_end_dropped(self): + bucket = _double_bucket(count=5) # samples at 0, P, 2P, 3P, 4P + trimmed = bc.trim_bucket(bucket, self._ns(PERIOD), self._ns(3 * PERIOD)) + self.assertEqual(bc.bucket_values(trimmed), [1.0, 2.0]) + self.assertEqual(bc.bucket_timestamps(trimmed), [T0_NANOS + PERIOD, T0_NANOS + 2 * PERIOD]) + + def test_sampling_clock_start_is_shifted_and_stays_a_clock(self): + trimmed = bc.trim_bucket(_double_bucket(count=5), self._ns(PERIOD + 1), self._ns(10 * PERIOD)) + self.assertEqual(trimmed.dataTimestamps.WhichOneof("value"), "samplingClock") + clock = trimmed.dataTimestamps.samplingClock + self.assertEqual(to_epoch_nanos(clock.startTime), T0_NANOS + 2 * PERIOD) + self.assertEqual(clock.periodNanos, PERIOD) + self.assertEqual(clock.count, 3) + + def test_trim_to_nothing_returns_none(self): + bucket = _double_bucket(count=3) + self.assertIsNone(bc.trim_bucket(bucket, self._ns(10 * PERIOD), self._ns(11 * PERIOD))) + self.assertIsNone(bc.trim_bucket(bucket, self._ns(1), self._ns(PERIOD))) # between two samples + + def test_range_covering_everything_returns_an_equal_copy(self): + bucket = _double_bucket(count=3) + trimmed = bc.trim_bucket(bucket, self._ns(-1), self._ns(10 * PERIOD)) + self.assertEqual(trimmed, bucket) + self.assertIsNot(trimmed, bucket) + + def test_input_is_not_modified_and_provider_and_metadata_carry_over(self): + metadata = dfb.column_metadata(tags=["t"]) + bucket = _bucket("PV:A", _clock(3), dfb.double_column("PV:A", [1.0, 2.0, 3.0], metadata=metadata)) + before = common_pb2.DataBucket() + before.CopyFrom(bucket) + trimmed = bc.trim_bucket(bucket, self._ns(PERIOD), self._ns(10 * PERIOD)) + self.assertEqual(bucket, before) + self.assertEqual((trimmed.pvName, trimmed.providerId, trimmed.providerName), ("PV:A", "prov-1", "daq")) + self.assertEqual(list(trimmed.dataValues.doubleColumn.metadata.tags), ["t"]) + + def test_timestamp_list_is_sliced(self): + nanos = [T0_NANOS, T0_NANOS + 10, T0_NANOS + 20, T0_NANOS + 30] + bucket = _bucket("PV:L", _raw_list(nanos), dfb.int32_column("PV:L", [1, 2, 3, 4])) + trimmed = bc.trim_bucket(bucket, self._ns(10), self._ns(30)) + self.assertEqual(bc.bucket_timestamps(trimmed), nanos[1:3]) + self.assertEqual(bc.bucket_values(trimmed), [2, 3]) + + def test_repeated_timestamps_at_each_bound(self): + # Two samples at begin, two at end: both copies at begin are kept, both at end dropped. + nanos = [T0_NANOS, T0_NANOS + 10, T0_NANOS + 10, T0_NANOS + 20, T0_NANOS + 30, T0_NANOS + 30] + bucket = _bucket("PV:D", _raw_list(nanos), dfb.double_column("PV:D", [0.0, 1.0, 2.0, 3.0, 4.0, 5.0])) + trimmed = bc.trim_bucket(bucket, self._ns(10), self._ns(30)) + self.assertEqual(bc.bucket_values(trimmed), [1.0, 2.0, 3.0]) + + def test_decreasing_axis_raises_naming_position(self): + nanos = [T0_NANOS, T0_NANOS + 20, T0_NANOS + 10] + bucket = _bucket("PV:R", _raw_list(nanos), dfb.double_column("PV:R", [1.0, 2.0, 3.0])) + with self.assertRaises(ValueError) as ctx: + bc.trim_bucket(bucket, self._ns(0), self._ns(100)) + self.assertIn("PV:R", str(ctx.exception)) + self.assertIn("position 2", str(ctx.exception)) + + def test_begin_not_before_end_raises(self): + bucket = _double_bucket() + for begin, end in ((10, 10), (20, 10)): + with self.subTest(begin=begin, end=end), self.assertRaises(ValueError) as ctx: + bc.trim_bucket(bucket, self._ns(begin), self._ns(end)) + self.assertIn("strictly before", str(ctx.exception)) + + def test_array_bucket_trims_in_whole_samples(self): + column = dfb.double_array_column("PV:W", [[[1.0, 2.0], [3.0, 4.0]], [[5.0, 6.0], [7.0, 8.0]], [[9.0] * 2] * 2]) + trimmed = bc.trim_bucket(_bucket("PV:W", _clock(3), column), self._ns(PERIOD), self._ns(2 * PERIOD)) + self.assertEqual(bc.bucket_values(trimmed), [[5.0, 6.0, 7.0, 8.0]]) + self.assertEqual(list(trimmed.dataValues.doubleArrayColumn.dimensions.dims), [2, 2]) + + def test_array_bucket_with_absent_dims_raises_rather_than_slicing(self): + column = common_pb2.DoubleArrayColumn(name="PV:W", values=[1.0, 2.0, 3.0, 4.0]) # no dims + bucket = _bucket("PV:W", _clock(2), column) + with self.assertRaises(ValueError) as ctx: + bc.trim_bucket(bucket, self._ns(0), self._ns(PERIOD)) + self.assertIn("dimensions", str(ctx.exception)) + + def test_image_bucket_keeps_its_descriptor(self): + column = dfb.image_column("PV:I", [b"a", b"b", b"c"], 4, 3, 1, "png") + trimmed = bc.trim_bucket(_bucket("PV:I", _clock(3), column), self._ns(PERIOD), self._ns(10 * PERIOD)) + self.assertEqual(bc.bucket_values(trimmed), [b"b", b"c"]) + self.assertEqual(trimmed.dataValues.imageColumn.imageDescriptor.encoding, "png") + + def test_serialized_bucket_raises(self): + bucket = _bucket("PV:S", _clock(2), dfb.serialized_column("PV:S", b"xy", "codec")) + with self.assertRaises(ValueError) as ctx: + bc.trim_bucket(bucket, self._ns(0), self._ns(PERIOD)) + self.assertIn("trimmed", str(ctx.exception)) + + def test_repeated_timestamps_trim_without_error(self): + nanos = [T0_NANOS, T0_NANOS, T0_NANOS] + bucket = _bucket("PV:D", _raw_list(nanos), dfb.double_column("PV:D", [1.0, 2.0, 3.0])) + self.assertEqual(bc.bucket_values(bc.trim_bucket(bucket, self._ns(0), self._ns(1))), [1.0, 2.0, 3.0]) + + +# ---------------------------------------------------------------------- +# grouping +# ---------------------------------------------------------------------- + + +class TestBucketsByPv(unittest.TestCase): + def test_interleaved_pvs_and_out_of_order_buckets(self): + a_late = _double_bucket("PV:A", start_nanos=T0_NANOS + 100 * PERIOD) + b_only = _double_bucket("PV:B") + a_early = _double_bucket("PV:A") + groups = bc.buckets_by_pv([a_late, b_only, a_early]) + self.assertEqual(list(groups), ["PV:A", "PV:B"]) + self.assertEqual(groups["PV:A"], [a_early, a_late]) + self.assertEqual(groups["PV:B"], [b_only]) + + def test_equal_first_timestamps_keep_input_order(self): + first = _double_bucket("PV:A", values=[1.0] * 5) + second = _double_bucket("PV:A", values=[2.0] * 5) + self.assertEqual(bc.buckets_by_pv([first, second])["PV:A"], [first, second]) + + def test_empty_input(self): + self.assertEqual(bc.buckets_by_pv([]), {}) + + +# ---------------------------------------------------------------------- +# pandas ([analysis] extra) +# ---------------------------------------------------------------------- + + +@unittest.skipUnless(_HAVE_ANALYSIS, "requires the [analysis] extra (pandas)") +class TestBucketsToDataFrames(unittest.TestCase): + def test_one_frame_per_pv_with_exact_index(self): + a1 = _double_bucket("PV:A", count=2) + a2 = _double_bucket("PV:A", count=2, start_nanos=T0_NANOS + 10 * PERIOD, values=[7.0, 8.0]) + b1 = _bucket("PV:B", _clock(3), dfb.int32_column("PV:B", [1, 2, 3])) + frames = bc.buckets_to_dataframes([a2, b1, a1]) + self.assertEqual(list(frames), ["PV:A", "PV:B"]) + df_a = frames["PV:A"] + self.assertEqual(list(df_a.columns), ["PV:A"]) + self.assertEqual(list(df_a["PV:A"]), [0.0, 1.0, 7.0, 8.0]) + self.assertEqual( + [ts.value for ts in df_a.index], + [T0_NANOS, T0_NANOS + PERIOD, T0_NANOS + 10 * PERIOD, T0_NANOS + 11 * PERIOD], + ) + self.assertEqual(str(df_a.index.tz), "UTC") + self.assertEqual(str(frames["PV:B"]["PV:B"].dtype), "int32") + + def test_narrow_dtypes_survive_concatenation(self): + f1 = _bucket("PV:F", _clock(2), dfb.float_column("PV:F", [0.5, 1.5])) + f2 = _bucket("PV:F", _clock(2, start_nanos=T0_NANOS + 10 * PERIOD), dfb.float_column("PV:F", [2.5, 3.5])) + self.assertEqual(str(bc.buckets_to_dataframes([f1, f2])["PV:F"]["PV:F"].dtype), "float32") + + def test_every_non_serialized_arm_converts(self): + for column, expected in _all_arm_cases(): + with self.subTest(kind=type(column).__name__): + df = bc.buckets_to_dataframes([_bucket("PV:X", _clock(3), column)])["PV:X"] + self.assertEqual(list(df["PV:X"]), expected) + + def test_overlapping_buckets_are_kept(self): + a1 = _double_bucket("PV:A", count=3) + a2 = _double_bucket("PV:A", count=3, start_nanos=T0_NANOS + PERIOD) + df = bc.buckets_to_dataframes([a1, a2])["PV:A"] + self.assertEqual(len(df), 6) + self.assertFalse(df.index.is_unique) + + def test_column_named_for_the_pv(self): + bucket = _bucket("PV:A", _clock(2), dfb.double_column("ingested-name", [1.0, 2.0])) + self.assertEqual(list(bc.buckets_to_dataframes([bucket])["PV:A"].columns), ["PV:A"]) + + def test_buckets_attr_descriptors(self): + a1 = _double_bucket("PV:A", count=2, provider_id="p1", provider_name="one") + a2 = _double_bucket("PV:A", count=3, start_nanos=T0_NANOS + 10 * PERIOD, provider_id="p2", provider_name="two") + entries = bc.buckets_to_dataframes([a1, a2])["PV:A"].attrs["buckets"] + self.assertEqual( + [ + (e["first_nanos"], e["last_nanos"], e["sample_count"], e["provider_id"], e["provider_name"]) + for e in entries + ], + [ + (T0_NANOS, T0_NANOS + PERIOD, 2, "p1", "one"), + (T0_NANOS + 10 * PERIOD, T0_NANOS + 12 * PERIOD, 3, "p2", "two"), + ], + ) + + def test_uniform_metadata_is_summarized(self): + metadata = dfb.column_metadata(tags=["vacuum"], attributes={"unit": "Torr"}) + buckets = [ + _bucket( + "PV:A", + _clock(2, start_nanos=T0_NANOS + i * 10 * PERIOD), + dfb.double_column("PV:A", [1.0, 2.0], metadata=metadata), + ) + for i in range(2) + ] + attrs = bc.buckets_to_dataframes(buckets)["PV:A"].attrs + self.assertEqual(attrs["column_metadata"]["PV:A"]["tags"], ["vacuum"]) + self.assertEqual(attrs["buckets"][1]["column_metadata"]["attributes"], {"unit": "Torr"}) + + def test_differing_metadata_is_kept_per_bucket_only(self): + buckets = [ + _bucket( + "PV:A", + _clock(2, start_nanos=T0_NANOS + i * 10 * PERIOD), + dfb.double_column("PV:A", [1.0, 2.0], metadata=dfb.column_metadata(tags=[f"run-{i}"])), + ) + for i in range(2) + ] + attrs = bc.buckets_to_dataframes(buckets)["PV:A"].attrs + self.assertNotIn("column_metadata", attrs) + self.assertEqual([e["column_metadata"]["tags"] for e in attrs["buckets"]], [["run-0"], ["run-1"]]) + + def test_exclude_column_metadata_keeps_structural_attrs(self): + bucket = _bucket( + "PV:E", _clock(3), dfb.enum_column("PV:E", [0, 1, 2], "mode:v1", metadata=dfb.column_metadata(tags=["t"])) + ) + attrs = bc.buckets_to_dataframes([bucket], exclude_column_metadata=True)["PV:E"].attrs + self.assertNotIn("column_metadata", attrs) + self.assertNotIn("column_metadata", attrs["buckets"][0]) + self.assertEqual(attrs["enum_ids"], {"PV:E": "mode:v1"}) + + def test_structural_attrs_for_arrays_images_structs(self): + array = _bucket("PV:W", _clock(2), dfb.double_array_column("PV:W", [[[1.0, 2.0]], [[3.0, 4.0]]])) + image = _bucket("PV:I", _clock(2), dfb.image_column("PV:I", [b"a", b"b"], 4, 3, 1, "png")) + struct = _bucket("PV:S", _clock(2), dfb.struct_column("PV:S", [b"x", b"y"], "schema:v1")) + frames = bc.buckets_to_dataframes([array, image, struct]) + self.assertEqual(frames["PV:W"].attrs["dimensions"], {"PV:W": [1, 2]}) + self.assertEqual(frames["PV:I"].attrs["image_descriptors"]["PV:I"]["encoding"], "png") + self.assertEqual(frames["PV:S"].attrs["schema_ids"], {"PV:S": "schema:v1"}) + + def test_kind_mismatch_raises(self): + a1 = _double_bucket("PV:A", count=2) + a2 = _bucket("PV:A", _clock(2, start_nanos=T0_NANOS + 10 * PERIOD), dfb.int32_column("PV:A", [1, 2])) + with self.assertRaises(ValueError) as ctx: + bc.buckets_to_dataframes([a1, a2]) + message = str(ctx.exception) + self.assertIn("PV:A", message) + self.assertIn(str(T0_NANOS), message) + self.assertIn(str(T0_NANOS + 10 * PERIOD), message) + + def test_structure_mismatch_raises(self): + cases = { + "enum": (dfb.enum_column("PV:A", [0, 1], "a:v1"), dfb.enum_column("PV:A", [0, 1], "b:v1")), + "dims": ( + dfb.double_array_column("PV:A", [[1.0, 2.0, 3.0, 4.0]] * 2), + dfb.double_array_column("PV:A", [[[1.0, 2.0], [3.0, 4.0]]] * 2), + ), + "image": ( + dfb.image_column("PV:A", [b"a", b"b"], 4, 3, 1, "png"), + dfb.image_column("PV:A", [b"a", b"b"], 4, 3, 1, "jpeg"), + ), + "struct": (dfb.struct_column("PV:A", [b"x"] * 2, "s:v1"), dfb.struct_column("PV:A", [b"x"] * 2, "s:v2")), + } + for label, (first, second) in cases.items(): + with self.subTest(label), self.assertRaises(ValueError): + bc.buckets_to_dataframes( + [ + _bucket("PV:A", _clock(2), first), + _bucket("PV:A", _clock(2, start_nanos=T0_NANOS + 10 * PERIOD), second), + ] + ) + + def test_serialized_bucket_refused(self): + bucket = _bucket("PV:S", _clock(2), dfb.serialized_column("PV:S", b"xy", "my-codec")) + with self.assertRaises(ValueError) as ctx: + bc.buckets_to_dataframes([bucket]) + self.assertIn("my-codec", str(ctx.exception)) + + def test_time_range_trims_and_drops_empty_pvs(self): + a = _double_bucket("PV:A", count=5) + b = _double_bucket("PV:B", count=2, start_nanos=T0_NANOS + 100 * PERIOD) + frames = bc.buckets_to_dataframes( + [a, b], time_range=(from_epoch_nanos(T0_NANOS + PERIOD), from_epoch_nanos(T0_NANOS + 3 * PERIOD)) + ) + self.assertEqual(list(frames), ["PV:A"]) + self.assertEqual(list(frames["PV:A"]["PV:A"]), [1.0, 2.0]) + self.assertEqual(frames["PV:A"].attrs["buckets"][0]["sample_count"], 2) + + def test_time_range_begin_not_before_end_raises(self): + same = from_epoch_nanos(T0_NANOS) + with self.assertRaises(ValueError) as ctx: + bc.buckets_to_dataframes([_double_bucket()], time_range=(same, same)) + self.assertIn("strictly before", str(ctx.exception)) + + def test_repeated_timestamps_convert(self): + nanos = [T0_NANOS, T0_NANOS, T0_NANOS + 1] + bucket = _bucket("PV:D", _raw_list(nanos), dfb.double_column("PV:D", [1.0, 2.0, 3.0])) + df = bc.buckets_to_dataframes([bucket])["PV:D"] + self.assertEqual([ts.value for ts in df.index], nanos) + + def test_empty_input(self): + self.assertEqual(bc.buckets_to_dataframes([]), {}) + + +@unittest.skipUnless(_HAVE_ANALYSIS, "requires the [analysis] extra (pandas)") +class TestQueryBucketsToDataFrames(unittest.TestCase): + def setUp(self): + self.params = QueryParams( + from_epoch_nanos(T0_NANOS + PERIOD), + from_epoch_nanos(T0_NANOS + 3 * PERIOD), + pv_selector=PvQuery.name_list(["PV:A"]), + ) + + def _pages(self, *bucket_lists): + pages = [] + for buckets in bucket_lists: + response = query_pb2.QueryBucketsResponse() + response.bucketQueryResult.dataBuckets.extend(buckets) + pages.append(QueryBucketsApiResult(is_error=False, message="", response=response)) + client = Mock() + client.iter_query_buckets = Mock(return_value=iter(pages)) + return client + + def test_assembles_a_pv_split_across_pages(self): + a1 = _double_bucket("PV:A", count=2) + a2 = _double_bucket("PV:A", count=2, start_nanos=T0_NANOS + 2 * PERIOD, values=[2.0, 3.0]) + frames = bc.query_buckets_to_dataframes(self._pages([a1], [a2]), self.params) + self.assertEqual(list(frames["PV:A"]["PV:A"]), [0.0, 1.0, 2.0, 3.0]) + + def test_trim_uses_the_params_range(self): + frames = bc.query_buckets_to_dataframes(self._pages([_double_bucket("PV:A")]), self.params, trim=True) + self.assertEqual(list(frames["PV:A"]["PV:A"]), [1.0, 2.0]) + + def test_untrimmed_by_default(self): + frames = bc.query_buckets_to_dataframes(self._pages([_double_bucket("PV:A")]), self.params) + self.assertEqual(len(frames["PV:A"]), 5) + + def test_max_buckets(self): + client = self._pages([_double_bucket("PV:A")], [_double_bucket("PV:B")]) + with self.assertRaises(ValueError) as ctx: + bc.query_buckets_to_dataframes(client, self.params, max_buckets=1) + self.assertIn("max_buckets=1", str(ctx.exception)) + client = self._pages([_double_bucket("PV:A")], [_double_bucket("PV:B")]) + self.assertEqual(len(bc.query_buckets_to_dataframes(client, self.params, max_buckets=2)), 2) + + def test_result_to_dataframes_delegates(self): + response = query_pb2.QueryBucketsResponse() + response.bucketQueryResult.dataBuckets.extend([_double_bucket("PV:A")]) + result = QueryBucketsApiResult(is_error=False, message="", response=response) + frames = result.to_dataframes(time_range=(from_epoch_nanos(T0_NANOS), from_epoch_nanos(T0_NANOS + PERIOD))) + self.assertEqual(list(frames["PV:A"]["PV:A"]), [0.0]) + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/unit/test_data_frame_conversions.py b/tests/unit/test_data_frame_conversions.py index a34c19d..289c726 100644 --- a/tests/unit/test_data_frame_conversions.py +++ b/tests/unit/test_data_frame_conversions.py @@ -59,7 +59,7 @@ def test_frame_without_axis_raises(self): class TestColumnValues(unittest.TestCase): - """column_values() is a standalone per-column converter, reusable by the bucket query (#16).""" + """column_values() is a standalone per-column converter, reused by bucket_conversions for the bucket query (#16).""" def test_each_scalar_column_type(self): cases = [ @@ -827,8 +827,8 @@ def test_other_column_kinds_have_neither(self): class TestBuilderReadBackSymmetry(unittest.TestCase): """ Each non-scalar builder (#17, D7) round-trips through data_frame() and the read side: values, one entry per - sample, plus the structural field the payload needs. Live read-back waits for the bucket query (#16), so this - is the check that the write and read halves agree on layout -- above all, row-major array flattening. + sample, plus the structural field the payload needs. This is the offline check that the write and read halves + agree on layout -- above all, row-major array flattening; live read-back goes through the bucket query (#16). """ def test_array_builders_round_trip_values_and_dims(self): diff --git a/tests/unit/test_query_client.py b/tests/unit/test_query_client.py index 72a5925..9b6f4d9 100644 --- a/tests/unit/test_query_client.py +++ b/tests/unit/test_query_client.py @@ -14,12 +14,13 @@ from dp_python_lib.client.query_client import ( ConfigQuery, PvQuery, + QueryBucketsApiResult, QueryClient, QueryParams, QuerySamplesApiResult, SampleStatusFilter, ) -from dp_python_lib.grpc import query_pb2 +from dp_python_lib.grpc import common_pb2, query_pb2 BEGIN = datetime(2024, 1, 1, tzinfo=timezone.utc) END = datetime(2024, 1, 2, tzinfo=timezone.utc) @@ -597,3 +598,265 @@ def test_helper_built_selectors_pass_validation(self): if __name__ == "__main__": unittest.main() + + +# ---------------------------------------------------------------------- +# queryBuckets / queryBucketsStream (#16) +# ---------------------------------------------------------------------- + + +def _bucket_response(pv_names=(), next_page_token=""): + """Build a real QueryBucketsResponse carrying one empty-bodied DataBucket per PV name and a nextPageToken.""" + response = query_pb2.QueryBucketsResponse() + for pv_name in pv_names: + response.bucketQueryResult.dataBuckets.add(pvName=pv_name) + response.bucketQueryResult.nextPageToken = next_page_token + return response + + +def _bucket_exceptional_response(message): + """Build a real QueryBucketsResponse carrying an exceptionalResult with the given message.""" + response = query_pb2.QueryBucketsResponse() + response.exceptionalResult.message = message + return response + + +class TestBuildQueryBucketsRequest(unittest.TestCase): + def setUp(self): + self.client = QueryClient(Mock()) + + def test_spec_matches_the_samples_builder(self): + params = QueryParams( + BEGIN, + END, + pv_selector=PvQuery.name_list(["A", "B"]), + config_criteria=[ConfigQuery.configuration_name(["cfg"])], + limit=50, + ) + buckets = self.client._build_query_buckets_request(params, page_token="tok") + samples = self.client._build_query_samples_request(params, page_token="tok") + self.assertIsInstance(buckets, query_pb2.QueryBucketsRequest) + self.assertEqual(buckets.querySpec, samples.querySpec) + self.assertEqual(buckets.executionOptions.limit, 50) + self.assertEqual(buckets.executionOptions.pageToken, "tok") + + def test_use_serialized_columns_is_forced_false(self): + params = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x")) + request = self.client._build_query_buckets_request(params) + self.assertFalse(request.resultRepresentation.useSerializedColumns) + + def test_exclude_column_metadata_is_passed_through(self): + for exclude in (False, True): + with self.subTest(exclude=exclude): + params = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x"), exclude_column_metadata=exclude) + request = self.client._build_query_buckets_request(params) + self.assertEqual(request.resultRepresentation.excludeColumnMetadata, exclude) + + def test_no_limit_no_token(self): + request = self.client._build_query_buckets_request(QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x"))) + self.assertEqual(request.executionOptions.limit, 0) + self.assertEqual(request.executionOptions.pageToken, "") + + def test_sample_status_filter_is_refused(self): + params = QueryParams( + BEGIN, END, pv_selector=PvQuery.pattern("x"), sample_status_filter=SampleStatusFilter.exclude("dq") + ) + with self.assertRaises(ValueError) as ctx: + self.client._build_query_buckets_request(params) + self.assertIn("sample_status_filter", str(ctx.exception)) + self.assertIn("query_samples()", str(ctx.exception)) + + def test_refusal_happens_before_any_rpc(self): + stub = Mock() + self.client._stub = stub + params = QueryParams( + BEGIN, END, pv_selector=PvQuery.pattern("x"), sample_status_filter=SampleStatusFilter.include("dq") + ) + with self.assertRaises(ValueError): + self.client.query_buckets(params) + with self.assertRaises(ValueError): + list(self.client.iter_query_buckets(params)) + with self.assertRaises(ValueError): + list(self.client.iter_query_buckets_stream(params)) + stub.queryBuckets.assert_not_called() + stub.queryBucketsStream.assert_not_called() + + +class TestQueryBucketsUnary(unittest.TestCase): + def setUp(self): + self.client = QueryClient(Mock()) + self.mock_stub = Mock() + self.client._stub = self.mock_stub + self.params = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x"), limit=7) + + def test_success(self): + self.mock_stub.queryBuckets.return_value = _bucket_response(["A", "B"], next_page_token="tok") + result = self.client.query_buckets(self.params, page_token="prev") + self.assertFalse(result.result_status.is_error) + self.assertEqual([b.pvName for b in result.data_buckets], ["A", "B"]) + self.assertIsInstance(result.data_buckets[0], common_pb2.DataBucket) + self.assertEqual(result.next_page_token, "tok") + sent = self.mock_stub.queryBuckets.call_args[0][0] + self.assertEqual(sent.executionOptions.pageToken, "prev") + self.assertEqual(sent.executionOptions.limit, 7) + + def test_empty_result_is_success(self): + self.mock_stub.queryBuckets.return_value = _bucket_response([]) + result = self.client.query_buckets(self.params) + self.assertFalse(result.result_status.is_error) + self.assertEqual(result.data_buckets, []) + self.assertEqual(result.next_page_token, "") + + def test_business_error(self): + self.mock_stub.queryBuckets.return_value = _bucket_exceptional_response("bad query") + result = self.client.query_buckets(self.params) + self.assertTrue(result.result_status.is_error) + self.assertEqual(result.result_status.message, "bad query") + self.assertEqual(result.data_buckets, []) + self.assertEqual(result.next_page_token, "") + + def test_unexpected_response_format(self): + self.mock_stub.queryBuckets.return_value = _response_with_field("somethingElse") + result = self.client.query_buckets(self.params) + self.assertTrue(result.result_status.is_error) + self.assertIn("neither exceptionalResult nor bucketQueryResult", result.result_status.message) + + def test_grpc_error(self): + err = grpc.RpcError() + err.details = lambda: "connection refused" + self.mock_stub.queryBuckets.side_effect = err + result = self.client.query_buckets(self.params) + self.assertTrue(result.result_status.is_error) + self.assertIn("gRPC error: connection refused", result.result_status.message) + + def test_unexpected_exception(self): + self.mock_stub.queryBuckets.side_effect = ValueError("boom") + result = self.client.query_buckets(self.params) + self.assertTrue(result.result_status.is_error) + self.assertIn("Unexpected error: boom", result.result_status.message) + + +class TestIterQueryBuckets(unittest.TestCase): + def setUp(self): + self.client = QueryClient(Mock()) + self.mock_stub = Mock() + self.client._stub = self.mock_stub + self.params = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x")) + + def test_paging_follows_tokens_and_stops_on_empty(self): + self.mock_stub.queryBuckets.side_effect = [ + _bucket_response(["A"], next_page_token="t1"), + _bucket_response(["A"], next_page_token="t2"), + _bucket_response(["B"], next_page_token=""), + ] + pages = list(self.client.iter_query_buckets(self.params)) + self.assertEqual([[b.pvName for b in p.data_buckets] for p in pages], [["A"], ["A"], ["B"]]) + sent_tokens = [c[0][0].executionOptions.pageToken for c in self.mock_stub.queryBuckets.call_args_list] + self.assertEqual(sent_tokens, ["", "t1", "t2"]) + + def test_error_page_raises_runtime_error(self): + self.mock_stub.queryBuckets.side_effect = [ + _bucket_response(["A"], next_page_token="t1"), + _bucket_exceptional_response("page boom"), + ] + collected = [] + with self.assertRaises(RuntimeError) as ctx: + for page in self.client.iter_query_buckets(self.params): + collected.append(page) + self.assertEqual(len(collected), 1) + self.assertIn("queryBuckets failed during paging: page boom", str(ctx.exception)) + + +class TestQueryBucketsStream(unittest.TestCase): + def setUp(self): + self.client = QueryClient(Mock()) + self.mock_stub = Mock() + self.client._stub = self.mock_stub + self.params = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("x")) + + def test_stream_yields_messages(self): + self.mock_stub.queryBucketsStream.return_value = iter([_bucket_response(["A"]), _bucket_response(["B"])]) + results = list(self.client.iter_query_buckets_stream(self.params)) + self.assertTrue(all(isinstance(r, QueryBucketsApiResult) for r in results)) + self.assertEqual([[b.pvName for b in r.data_buckets] for r in results], [["A"], ["B"]]) + + def test_stream_never_sends_page_token(self): + self.mock_stub.queryBucketsStream.return_value = iter([_bucket_response()]) + list(self.client.iter_query_buckets_stream(self.params)) + sent = self.mock_stub.queryBucketsStream.call_args[0][0] + self.assertIsInstance(sent, query_pb2.QueryBucketsRequest) + self.assertEqual(sent.executionOptions.pageToken, "") + + def test_empty_result_is_one_empty_message(self): + self.mock_stub.queryBucketsStream.return_value = iter([_bucket_response([])]) + results = list(self.client.iter_query_buckets_stream(self.params)) + self.assertEqual(len(results), 1) + self.assertEqual(results[0].data_buckets, []) + + def test_business_error_mid_stream_raises(self): + self.mock_stub.queryBucketsStream.return_value = iter( + [_bucket_response(["A"]), _bucket_exceptional_response("single bucket for pv A exceeds the limit")] + ) + collected = [] + with self.assertRaises(RuntimeError) as ctx: + for r in self.client.iter_query_buckets_stream(self.params): + collected.append(r) + self.assertEqual(len(collected), 1) + self.assertIn("queryBucketsStream failed during streaming", str(ctx.exception)) + self.assertIn("exceeds the limit", str(ctx.exception)) + + def test_grpc_error_mid_stream_raises(self): + err = grpc.RpcError() + err.details = lambda: "stream reset" + + def exploding_stream(): + yield _bucket_response(["A"]) + raise err + + self.mock_stub.queryBucketsStream.return_value = exploding_stream() + with self.assertRaises(RuntimeError) as ctx: + list(self.client.iter_query_buckets_stream(self.params)) + self.assertIn("gRPC error: stream reset", str(ctx.exception)) + + def test_unrecognized_message_raises(self): + self.mock_stub.queryBucketsStream.return_value = iter([_response_with_field("somethingElse")]) + with self.assertRaises(RuntimeError) as ctx: + list(self.client.iter_query_buckets_stream(self.params)) + self.assertIn("neither exceptionalResult nor bucketQueryResult", str(ctx.exception)) + + +class TestSharedStreamSender(unittest.TestCase): + """The samples stream sender now goes through _iter_stream(); its yielded error messages are unchanged.""" + + def setUp(self): + self.client = QueryClient(Mock()) + self.mock_stub = Mock() + self.client._stub = self.mock_stub + self.request = query_pb2.QuerySamplesRequest() + + def test_samples_error_texts(self): + err = grpc.RpcError() + err.details = lambda: "reset" + + def stream(): + yield _exceptional_response("biz") + yield _response_with_field("somethingElse") + raise err + + self.mock_stub.querySamplesStream.return_value = stream() + messages = [r.result_status.message for r in self.client._send_query_samples_stream(self.request)] + self.assertEqual( + messages, + [ + "biz", + "Unexpected response format: neither exceptionalResult nor sampleQueryResult found", + "gRPC error: reset", + ], + ) + + def test_unexpected_exception_text(self): + self.mock_stub.querySamplesStream.side_effect = ValueError("boom") + results = list(self.client._send_query_samples_stream(self.request)) + self.assertEqual(len(results), 1) + self.assertIsInstance(results[0], QuerySamplesApiResult) + self.assertEqual(results[0].result_status.message, "Unexpected error: boom") From c210c3697e91c967b65b7af2511d32543ed4f531 Mon Sep 17 00:00:00 2001 From: Craig McChesney Date: Thu, 1 Oct 2026 11:24:11 -0600 Subject: [PATCH 2/3] fix: only refuse serialized buckets that overlap time_range (#16 PR A review) buckets_to_dataframes() refused a SerializedDataColumn bucket before trimming, across every PV, so one serialized bucket anywhere in a result -- even wholly outside time_range -- blocked the conversion of every other PV. With time_range, buckets with no sample in the range are now dropped first (they would be trimmed away anyway), and only an overlapping serialized bucket is refused. Without time_range nothing changes: every bucket is kept, so a serialized one still raises. Also documents that a legacy DataColumn bypasses the consistency check (it has no column-level type), so its buckets widen through pandas rather than raising, and pins that with a test. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01CajTMkjkkzeoWXLSgMpn5k --- .../client/bucket_conversions.py | 24 +++++++++++-- tests/unit/test_bucket_conversions.py | 34 +++++++++++++++++++ 2 files changed, 55 insertions(+), 3 deletions(-) diff --git a/src/dp_python_lib/client/bucket_conversions.py b/src/dp_python_lib/client/bucket_conversions.py index 0132c5a..fe7225e 100644 --- a/src/dp_python_lib/client/bucket_conversions.py +++ b/src/dp_python_lib/client/bucket_conversions.py @@ -187,6 +187,11 @@ def _first_out_of_order(epoch_nanos: list[int]) -> int | None: return None +def _has_samples_in(bucket: common_pb2.DataBucket, begin_nanos: int, end_nanos: int) -> bool: + """Whether any of a bucket's timestamps falls in [begin, end). Reads only the axis, so any column kind works.""" + return any(begin_nanos <= nanos < end_nanos for nanos in bucket_timestamps(bucket)) + + def _trim_nanos(bucket: common_pb2.DataBucket, begin_nanos: int, end_nanos: int) -> common_pb2.DataBucket | None: """trim_bucket() over an already-validated range in epoch nanoseconds.""" column = bucket_column(bucket) @@ -301,6 +306,7 @@ def _structure(column: Any) -> tuple: def _check_consistent(pv_name: str, buckets: list[common_pb2.DataBucket]) -> None: """ Rejects a PV whose buckets differ in column kind or structure, which cannot be concatenated into one column. + A legacy DataColumn carries no column-level type, so its buckets always pass (see buckets_to_dataframes()). Concatenating anyway would silently widen dtypes (double then int32) or mix payloads that mean different things (two enumerations, two array shapes). @@ -397,6 +403,12 @@ def buckets_to_dataframes( them. These are structural -- the values cannot be interpreted without them -- so they are present even under exclude_column_metadata. + A legacy DataColumn is the exception to the consistency check: it has no column-level type, so its buckets are + always treated as one structure even when their DataValues use different arms. Their parts then concatenate + with pandas' own widening -- integers with gaps, or ints in one bucket and doubles in another, become float64, + and a mix with strings becomes object -- rather than raising. Read such buckets with bucket_values() when the + per-sample arm matters. + :param buckets: The buckets, in any order (e.g. the data_buckets of every page or stream message). :param time_range: Optional (begin, end) to trim each bucket to [begin, end) exactly (see trim_bucket()); None (the default) keeps every bucket whole, as the server returned it. A PV with no sample left in the range is @@ -406,7 +418,8 @@ def buckets_to_dataframes( :raises ImportError: if the [analysis] extra is not installed. :raises ValueError: if time_range's begin is not before its end; if a PV's buckets differ in column kind or structure; or if a bucket is malformed, holds a SerializedDataColumn, or (when trimming) has a decreasing - time axis. + time axis. With time_range, only buckets holding a sample in the range are checked: one lying wholly + outside it is dropped, a serialized one included. """ pd = _require_pandas() @@ -414,12 +427,17 @@ def buckets_to_dataframes( frames: dict[str, Any] = {} for pv_name, group in buckets_by_pv(buckets).items(): + if range_nanos is not None: + # A serialized bucket with no sample in the range would be dropped by the trim like any other, so it is + # refused only when it overlaps. Refusing it up front made one serialized PV anywhere in a result -- + # even entirely outside the range -- block the conversion of every PV. + group = [bucket for bucket in group if _has_samples_in(bucket, *range_nanos)] + if not group: + continue for bucket in group: _refuse_serialized(bucket, bucket_column(bucket), "converted to a pandas DataFrame") if range_nanos is not None: group = [trimmed for bucket in group if (trimmed := _trim_nanos(bucket, *range_nanos)) is not None] - if not group: - continue _check_consistent(pv_name, group) frames[pv_name] = _pv_dataframe(pd, pv_name, group, exclude_column_metadata) return frames diff --git a/tests/unit/test_bucket_conversions.py b/tests/unit/test_bucket_conversions.py index 15c81aa..975a810 100644 --- a/tests/unit/test_bucket_conversions.py +++ b/tests/unit/test_bucket_conversions.py @@ -463,6 +463,40 @@ def test_serialized_bucket_refused(self): bc.buckets_to_dataframes([bucket]) self.assertIn("my-codec", str(ctx.exception)) + def test_serialized_bucket_outside_time_range_does_not_block_other_pvs(self): + a = _double_bucket("PV:A", count=5) + s = _bucket("PV:S", _clock(2, start_nanos=T0_NANOS + 100 * PERIOD), dfb.serialized_column("PV:S", b"xy", "c")) + frames = bc.buckets_to_dataframes( + [a, s], time_range=(from_epoch_nanos(T0_NANOS), from_epoch_nanos(T0_NANOS + 5 * PERIOD)) + ) + self.assertEqual(list(frames), ["PV:A"]) + + def test_serialized_bucket_outside_time_range_dropped_from_its_own_pv(self): + whole = _double_bucket("PV:S", count=2) + s = _bucket("PV:S", _clock(2, start_nanos=T0_NANOS + 100 * PERIOD), dfb.serialized_column("PV:S", b"xy", "c")) + frames = bc.buckets_to_dataframes( + [whole, s], time_range=(from_epoch_nanos(T0_NANOS), from_epoch_nanos(T0_NANOS + 5 * PERIOD)) + ) + self.assertEqual(list(frames["PV:S"]["PV:S"]), [0.0, 1.0]) + + def test_serialized_bucket_overlapping_time_range_still_refused(self): + a = _double_bucket("PV:A", count=5) + s = _bucket("PV:S", _clock(2), dfb.serialized_column("PV:S", b"xy", "my-codec")) + with self.assertRaises(ValueError) as ctx: + bc.buckets_to_dataframes( + [a, s], time_range=(from_epoch_nanos(T0_NANOS), from_epoch_nanos(T0_NANOS + 5 * PERIOD)) + ) + self.assertIn("my-codec", str(ctx.exception)) + + def test_legacy_data_column_buckets_widen_rather_than_raise(self): + # Documented exception to the consistency check: a DataColumn has no column-level type. + ints = _bucket("PV:L", _clock(2), dfb.data_column("PV:L", [1, None])) + doubles = _bucket("PV:L", _clock(2, start_nanos=T0_NANOS + 10 * PERIOD), dfb.data_column("PV:L", [2.5, 3.5])) + df = bc.buckets_to_dataframes([ints, doubles])["PV:L"] + self.assertEqual(str(df["PV:L"].dtype), "float64") + self.assertEqual(df["PV:L"].tolist()[0], 1.0) + self.assertEqual(df["PV:L"].tolist()[2:], [2.5, 3.5]) + def test_time_range_trims_and_drops_empty_pvs(self): a = _double_bucket("PV:A", count=5) b = _double_bucket("PV:B", count=2, start_nanos=T0_NANOS + 100 * PERIOD) From eec8e0e8b4255ee46830aedb4bd839f2634afab7 Mon Sep 17 00:00:00 2001 From: Craig McChesney Date: Thu, 1 Oct 2026 11:29:15 -0600 Subject: [PATCH 3/3] fix: honest absent metadata and clearer errors in bucket_conversions (#16 PR A review) - A bucket carrying no ColumnMetadata (as after a query that set excludeColumnMetadata, which a page's to_dataframes() cannot know) is now reported as column_metadata=None in attrs["buckets"], not as column_metadata_dict()'s empty containers, which read as metadata the server returned. The per-PV summary is set only when every bucket carries metadata and it all agrees. - attrs["buckets"] records each bucket's stored column_name, since the frame's column is always renamed to the PV. - trim_bucket() and time_range= errors name the caller's own parameter instead of a "time_range" trim_bucket() was never passed. - The kind/structure mismatch error renders readably, e.g. "EnumColumn (enumId 'a:v1')", instead of a raw tuple. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01CajTMkjkkzeoWXLSgMpn5k --- CLAUDE.md | 2 +- .../client/bucket_conversions.py | 70 +++++++++++++------ src/dp_python_lib/client/query_client.py | 4 +- tests/unit/test_bucket_conversions.py | 38 ++++++++++ 4 files changed, 92 insertions(+), 22 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index ba651ac..e20cd87 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -245,7 +245,7 @@ plan documents one change, `CLAUDE.md` documents the invariant it established. - `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, metadata), `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/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 diff --git a/src/dp_python_lib/client/bucket_conversions.py b/src/dp_python_lib/client/bucket_conversions.py index fe7225e..44c295a 100644 --- a/src/dp_python_lib/client/bucket_conversions.py +++ b/src/dp_python_lib/client/bucket_conversions.py @@ -28,12 +28,11 @@ from dp_python_lib.client import data_frame_conversions as dfc -# _slice_frame() and _time_range() are private to data_frame.py, which is the shared, kind-neutral home of the -# DataFrame builders rather than a feature module. Reusing them keeps one implementation of the integer-nanosecond -# SamplingClock shift and of the "begin strictly before end" rule and its message. -from dp_python_lib.client.data_frame import _COLUMN_FIELD_BY_TYPE, _slice_frame, _time_range +# _slice_frame() is private to data_frame.py, which is the shared, kind-neutral home of the DataFrame builders rather +# than a feature module. Reusing it keeps one implementation of the integer-nanosecond SamplingClock shift. +from dp_python_lib.client.data_frame import _COLUMN_FIELD_BY_TYPE, _slice_frame from dp_python_lib.client.sample_status_conversions import expand_data_timestamps -from dp_python_lib.client.time_conversions import TimestampInput, to_epoch_nanos +from dp_python_lib.client.time_conversions import TimestampInput, to_epoch_nanos, to_timestamp from dp_python_lib.grpc import common_pb2 @@ -166,17 +165,23 @@ def bucket_values(bucket: common_pb2.DataBucket) -> list: return _aligned_values(bucket, column, len(bucket_timestamps(bucket))) -def _range_nanos(begin: TimestampInput, end: TimestampInput) -> tuple[int, int]: +def _range_nanos(begin: TimestampInput, end: TimestampInput, caller: str) -> tuple[int, int]: """ Converts a (begin, end) pair to epoch nanoseconds, requiring begin strictly before end. A reversed or empty range is rejected rather than trimming to nothing, which would be indistinguishable from a - valid range that simply holds no samples. + valid range that simply holds no samples. The rule and the message's shape are data_frame._time_range()'s, but + the message names the caller's own parameter (trim_bucket()'s begin/end, or time_range=). + :param caller: How the range was passed, e.g. "trim_bucket()" or "time_range"; starts the error message. :raises ValueError: if begin is not strictly before end, or either bound is not a supported time input. """ - time_range = _time_range((begin, end)) - return to_epoch_nanos(time_range.beginTime), to_epoch_nanos(time_range.endTime) + begin_nanos, end_nanos = to_epoch_nanos(to_timestamp(begin)), to_epoch_nanos(to_timestamp(end)) + if begin_nanos >= end_nanos: + raise ValueError( + f"{caller} requires begin strictly before end; got begin {begin_nanos} ns and end {end_nanos} ns" + ) + return begin_nanos, end_nanos def _first_out_of_order(epoch_nanos: list[int]) -> int | None: @@ -247,7 +252,7 @@ def trim_bucket( :raises ValueError: if begin is not strictly before end; if the bucket is malformed, holds a SerializedDataColumn (an opaque payload cannot be sliced), is not aligned with its axis, or has a decreasing time axis. """ - begin_nanos, end_nanos = _range_nanos(begin, end) + begin_nanos, end_nanos = _range_nanos(begin, end, "trim_bucket()") return _trim_nanos(bucket, begin_nanos, end_nanos) @@ -303,6 +308,21 @@ def _structure(column: Any) -> tuple: ) +def _describe_structure(structure: tuple) -> str: + """Renders a _structure() tuple for an error message, e.g. "EnumColumn (enumId 'mode:v1')".""" + kind, enum_id, dims, descriptor, schema_id = structure + details = [] + if enum_id is not None: + details.append(f"enumId {enum_id!r}") + if dims is not None: + details.append(f"dims {list(dims)}") + if descriptor is not None: + details.append("image " + ", ".join(f"{key}={value!r}" for key, value in descriptor)) + if schema_id is not None: + details.append(f"schemaId {schema_id!r}") + return f"{kind} ({'; '.join(details)})" if details else kind + + def _check_consistent(pv_name: str, buckets: list[common_pb2.DataBucket]) -> None: """ Rejects a PV whose buckets differ in column kind or structure, which cannot be concatenated into one column. @@ -320,23 +340,29 @@ def _check_consistent(pv_name: str, buckets: list[common_pb2.DataBucket]) -> Non if actual != expected: raise ValueError( f"PV '{pv_name}' has buckets with different column kinds or structure (the bucket starting at " - f"{_first_nanos(first)} ns has {expected}, the one starting at {_first_nanos(bucket)} ns has " - f"{actual}), so they cannot form one column; read them separately with bucket_to_data_frame() or " - f"bucket_values()" + f"{_first_nanos(first)} ns has {_describe_structure(expected)}, the one starting at " + f"{_first_nanos(bucket)} ns has {_describe_structure(actual)}), so they cannot form one column; " + f"read them separately with bucket_to_data_frame() or bucket_values()" ) def _bucket_descriptor(bucket: common_pb2.DataBucket, epoch_nanos: list[int], exclude_column_metadata: bool) -> dict: """The per-bucket entry of df.attrs["buckets"].""" + column = bucket_column(bucket) descriptor: dict[str, Any] = { "first_nanos": epoch_nanos[0], "last_nanos": epoch_nanos[-1], "sample_count": len(epoch_nanos), "provider_id": bucket.providerId, "provider_name": bucket.providerName, + # The frame's column is named for the PV; this is the name the column was actually stored under. + "column_name": column.name, } if not exclude_column_metadata: - descriptor["column_metadata"] = dfc.column_metadata_dict(bucket_column(bucket)) + # None, not column_metadata_dict()'s empty containers, when the bucket carries no metadata at all -- as + # when the query itself set excludeColumnMetadata, which a page's to_dataframes() cannot know. An empty + # summary would read as metadata the server returned. + descriptor["column_metadata"] = dfc.column_metadata_dict(column) if column.HasField("metadata") else None return descriptor @@ -347,8 +373,9 @@ def _pv_dataframe(pd: Any, pv_name: str, buckets: list[common_pb2.DataBucket], e for bucket in buckets: frame = bucket_to_data_frame(bucket) part = dfc.data_frame_to_pandas(frame, exclude_column_metadata=True) - # The column is named for the ingested column, which ingestion makes the PV name; name it for the PV - # regardless, so every part concatenates into the one column. + # Ingestion names a bucket's column for its PV, but the stored name is not guaranteed to match; name it for + # the PV regardless, so every part concatenates into the one column. The stored name stays in + # attrs["buckets"][i]["column_name"]. part.columns = [pv_name] parts.append(part) descriptors.append(_bucket_descriptor(bucket, dfc.data_frame_timestamps(frame), exclude_column_metadata)) @@ -374,7 +401,7 @@ def _pv_dataframe(pd: Any, pv_name: str, buckets: list[common_pb2.DataBucket], e # Each ingest carries its own metadata, so a PV's buckets can legitimately differ. The per-PV summary is set # only when they agree; the per-bucket copies in attrs["buckets"] are always there. first_metadata = descriptors[0]["column_metadata"] - if all(entry["column_metadata"] == first_metadata for entry in descriptors[1:]): + if first_metadata is not None and all(entry["column_metadata"] == first_metadata for entry in descriptors[1:]): attrs["column_metadata"] = {pv_name: first_metadata} df.attrs = attrs return df @@ -397,8 +424,11 @@ def buckets_to_dataframes( df.attrs carries: - "buckets": one dict per bucket in the frame, with 'first_nanos', 'last_nanos', 'sample_count', - 'provider_id', 'provider_name', and (unless exclude_column_metadata) 'column_metadata'. - - "column_metadata": {pv: dict}, only when every bucket's metadata is identical and not excluded. + 'provider_id', 'provider_name', 'column_name' (the name the column was stored under; the frame's column is + always named for the PV), and (unless exclude_column_metadata) 'column_metadata': a dict, or None when the + bucket carries no metadata -- e.g. because the query set exclude_column_metadata. + - "column_metadata": {pv: dict}, only when every bucket carries metadata, all identical, and it is not + excluded. - "enum_ids", "dimensions", "image_descriptors", "schema_ids": {pv: value} for the column kinds that carry them. These are structural -- the values cannot be interpreted without them -- so they are present even under exclude_column_metadata. @@ -423,7 +453,7 @@ def buckets_to_dataframes( """ pd = _require_pandas() - range_nanos = _range_nanos(*time_range) if time_range is not None else None + range_nanos = _range_nanos(*time_range, "time_range") if time_range is not None else None frames: dict[str, Any] = {} for pv_name, group in buckets_by_pv(buckets).items(): diff --git a/src/dp_python_lib/client/query_client.py b/src/dp_python_lib/client/query_client.py index ef85d74..725bf85 100644 --- a/src/dp_python_lib/client/query_client.py +++ b/src/dp_python_lib/client/query_client.py @@ -588,7 +588,9 @@ def to_dataframes( :param time_range: Optional (begin, end) to trim each bucket to [begin, end) exactly; None (the default) leaves the buckets whole, as the server returned them. - :param exclude_column_metadata: If True, do not attach per-bucket ColumnMetadata to the DataFrames. + :param exclude_column_metadata: If True, do not attach per-bucket ColumnMetadata to the DataFrames. When + the query itself excluded metadata, leaving this False reports each bucket's metadata as None (absent), + not as an empty summary. :return: A dict of PV name to pandas.DataFrame. """ from dp_python_lib.client import bucket_conversions diff --git a/tests/unit/test_bucket_conversions.py b/tests/unit/test_bucket_conversions.py index 975a810..b448848 100644 --- a/tests/unit/test_bucket_conversions.py +++ b/tests/unit/test_bucket_conversions.py @@ -264,6 +264,8 @@ def test_begin_not_before_end_raises(self): with self.subTest(begin=begin, end=end), self.assertRaises(ValueError) as ctx: bc.trim_bucket(bucket, self._ns(begin), self._ns(end)) self.assertIn("strictly before", str(ctx.exception)) + # Names trim_bucket()'s own parameters, not a time_range it was never passed. + self.assertTrue(str(ctx.exception).startswith("trim_bucket() requires"), str(ctx.exception)) def test_array_bucket_trims_in_whole_samples(self): column = dfb.double_array_column("PV:W", [[[1.0, 2.0], [3.0, 4.0]], [[5.0, 6.0], [7.0, 8.0]], [[9.0] * 2] * 2]) @@ -407,6 +409,30 @@ def test_differing_metadata_is_kept_per_bucket_only(self): self.assertNotIn("column_metadata", attrs) self.assertEqual([e["column_metadata"]["tags"] for e in attrs["buckets"]], [["run-0"], ["run-1"]]) + def test_absent_metadata_is_none_not_an_empty_summary(self): + # What a bucket looks like after the query set excludeColumnMetadata: no metadata field at all. + buckets = [_double_bucket("PV:A", count=2), _double_bucket("PV:A", count=2, start_nanos=T0_NANOS + 10 * PERIOD)] + self.assertFalse(bc.bucket_column(buckets[0]).HasField("metadata")) + attrs = bc.buckets_to_dataframes(buckets)["PV:A"].attrs + self.assertNotIn("column_metadata", attrs) + self.assertEqual([entry["column_metadata"] for entry in attrs["buckets"]], [None, None]) + + def test_mixed_present_and_absent_metadata_is_kept_per_bucket_only(self): + with_metadata = _bucket( + "PV:A", _clock(2), dfb.double_column("PV:A", [1.0, 2.0], metadata=dfb.column_metadata(tags=["t"])) + ) + without = _double_bucket("PV:A", count=2, start_nanos=T0_NANOS + 10 * PERIOD) + attrs = bc.buckets_to_dataframes([with_metadata, without])["PV:A"].attrs + self.assertNotIn("column_metadata", attrs) + self.assertEqual(attrs["buckets"][0]["column_metadata"]["tags"], ["t"]) + self.assertIsNone(attrs["buckets"][1]["column_metadata"]) + + def test_stored_column_name_is_recorded(self): + bucket = _bucket("PV:A", _clock(2), dfb.double_column("stored-name", [1.0, 2.0])) + df = bc.buckets_to_dataframes([bucket])["PV:A"] + self.assertEqual(list(df.columns), ["PV:A"]) + self.assertEqual(df.attrs["buckets"][0]["column_name"], "stored-name") + def test_exclude_column_metadata_keeps_structural_attrs(self): bucket = _bucket( "PV:E", _clock(3), dfb.enum_column("PV:E", [0, 1, 2], "mode:v1", metadata=dfb.column_metadata(tags=["t"])) @@ -434,6 +460,17 @@ def test_kind_mismatch_raises(self): self.assertIn("PV:A", message) self.assertIn(str(T0_NANOS), message) self.assertIn(str(T0_NANOS + 10 * PERIOD), message) + self.assertIn("has DoubleColumn,", message) + self.assertIn("has Int32Column)", message) + self.assertNotIn("None", message) + + def test_structure_mismatch_message_names_the_field(self): + a1 = _bucket("PV:A", _clock(2), dfb.enum_column("PV:A", [0, 1], "a:v1")) + a2 = _bucket("PV:A", _clock(2, start_nanos=T0_NANOS + 10 * PERIOD), dfb.enum_column("PV:A", [0, 1], "b:v1")) + with self.assertRaises(ValueError) as ctx: + bc.buckets_to_dataframes([a1, a2]) + self.assertIn("EnumColumn (enumId 'a:v1')", str(ctx.exception)) + self.assertIn("EnumColumn (enumId 'b:v1')", str(ctx.exception)) def test_structure_mismatch_raises(self): cases = { @@ -512,6 +549,7 @@ def test_time_range_begin_not_before_end_raises(self): with self.assertRaises(ValueError) as ctx: bc.buckets_to_dataframes([_double_bucket()], time_range=(same, same)) self.assertIn("strictly before", str(ctx.exception)) + self.assertTrue(str(ctx.exception).startswith("time_range requires"), str(ctx.exception)) def test_repeated_timestamps_convert(self): nanos = [T0_NANOS, T0_NANOS, T0_NANOS + 1]