From 00d94877bcf6456ffc755735e14c576df404060e Mon Sep 17 00:00:00 2001 From: Craig McChesney Date: Thu, 1 Oct 2026 11:41:27 -0600 Subject: [PATCH 1/2] Live bucket query tests and docs (#16 PR B) - tests/integration/test_query_buckets_integration.py: a SamplingClock bucket returned verbatim, whole boundary buckets and exact trimming, limit counting buckets across pages, the stream matching the unary pages, an empty result, and ColumnMetadata provenance read back live and excluded. - test_non_scalar_columns_reach_success becomes test_non_scalar_columns_read_back_exactly: values, dims, image descriptor, schemaId, and the serialized payload and encoding. - Cookbook: a "Whole buckets" section in query.md; ingestion.md reads its array, image, and struct columns back (and its imports block gains what that needs); #16 pointers in conventions.md and sample-status.md. - README, NEXT.md, and a CLAUDE.md "Bucket Query API" section recording the server facts and design rules as invariants. The new and changed recipes were run as one script against a live stack. Closes #16 Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01CajTMkjkkzeoWXLSgMpn5k --- .dev/tools/check-cookbook-snippets.py | 1 + CLAUDE.md | 79 +++++- README.md | 10 +- doc/cookbook/README.md | 2 +- doc/cookbook/conventions.md | 4 +- doc/cookbook/ingestion.md | 28 +- doc/cookbook/query.md | 151 +++++++++- doc/cookbook/sample-status.md | 7 +- doc/release-notes/NEXT.md | 11 +- .../test_ingestion_client_integration.py | 52 +++- .../test_query_buckets_integration.py | 267 ++++++++++++++++++ 11 files changed, 580 insertions(+), 32 deletions(-) create mode 100644 tests/integration/test_query_buckets_integration.py diff --git a/.dev/tools/check-cookbook-snippets.py b/.dev/tools/check-cookbook-snippets.py index f2e9a13..4c692dd 100755 --- a/.dev/tools/check-cookbook-snippets.py +++ b/.dev/tools/check-cookbook-snippets.py @@ -132,6 +132,7 @@ from dp_python_lib.client import sample_status_conversions as ssc from dp_python_lib.client import data_frame as dfb from dp_python_lib.client import data_frame_conversions as dfc +from dp_python_lib.client import bucket_conversions as bc client: MldpClient = MldpClient() # client.annotation and client.query are typed `X | None`, which is honest: they are None when MldpClient is given diff --git a/CLAUDE.md b/CLAUDE.md index e20cd87..d8d6f72 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -262,7 +262,8 @@ plan documents one change, `CLAUDE.md` documents the invariant it established. - `tests/unit/test_export_client.py` - Unit tests for ExportClient (`ExportFormat` mapping and unreachable `UNSPECIFIED`, `calculations_spec()`, the zero-source rejection, three-tier error handling) - `tests/unit/test_annotation_client.py` - Unit tests pinning the `AnnotationClient` facade wiring (every feature client present, one shared channel, one stub apiece) - `tests/integration/ingest_support.py` - The integration tests' one way to put samples in the archive: `register_provider()` and `ingest_confirmed()`, which ingests a frame through `IngestionClient` and returns only once its request-status document says SUCCESS (raising otherwise). Every test needing archived data uses it, so none polls for bucket visibility: SUCCESS is written after the buckets. A setUpClass that cannot ingest turns the `RuntimeError`/`TimeoutError` into a skip -- `tests/integration/test_ingestion_client_integration.py` - Live ingestion (#17 PR B), each outcome confirmed through `await_request_statuses()`: unary ack + SUCCESS + exact read-back; an unknown `providerId` and a re-ingest of the same PV + first timestamp both **acked, then ERROR**; a client-streaming call with one server-rejected request (`rejected_request_ids`, the other two SUCCESS, the reject's status REJECTED); the bidi stream's in-order per-request results; a `split_data_frame()`-chunked frame reading back whole; array/image/struct/serialized columns reaching SUCCESS; and an oversized request failing with the `split_data_frame()` hint, unary and streamed. The reject comes from a frame spanning more than the server's one-day bucket-span cap, which the client deliberately leaves server-side +- `tests/integration/test_ingestion_client_integration.py` - Live ingestion (#17 PR B), each outcome confirmed through `await_request_statuses()`: unary ack + SUCCESS + exact read-back; an unknown `providerId` and a re-ingest of the same PV + first timestamp both **acked, then ERROR**; a client-streaming call with one server-rejected request (`rejected_request_ids`, the other two SUCCESS, the reject's status REJECTED); the bidi stream's in-order per-request results; a `split_data_frame()`-chunked frame reading back whole; array/image/struct/serialized columns read back exactly through the bucket query (#16); and an oversized request failing with the `split_data_frame()` hint, unary and streamed. The reject comes from a frame spanning more than the server's one-day bucket-span cap, which the client deliberately leaves server-side +- `tests/integration/test_query_buckets_integration.py` - Live bucket query (#16 PR B) over data it ingests itself, one request per bucket: a `SamplingClock` bucket returned verbatim (axis, provider, column type); a sub-window returning the whole boundary bucket, then `trim_bucket()` and `query_buckets_to_dataframes(trim=True)` reducing it exactly; `limit` counting buckets, with two PVs reassembling across pages (and the server's observed `(pvName, firstTime)` order, which the client does not rely on but the cookbook describes); `max_buckets`; the stream returning the same buckets as the unary pages; an empty result as success; and `ColumnMetadata` provenance read back live, absent under `exclude_column_metadata` (and `None`, not an empty summary, in `attrs["buckets"]`) - `tests/integration/test_query_client_integration.py` - Live v2 query mechanics against whatever data exists, plus `TestQueryClosedLoop`, which ingests two PVs on different clocks and asserts exact values and timestamps, half-open trimming at both bounds, dense alignment (the slower PV's missing rows are unset `DataValue`s, not zeros), and in-order paging at a small `limit` - `tests/integration/test_datasets_annotations_integration.py` - Live-server round trip for datasets/annotations/calculations; ingests its own samples first, because `saveDataSet` requires archived PVs - `tests/integration/test_query_helper_relaxations_integration.py` - Live-server coverage for the #40 key-only `attributes()` search and the #41 browse-all `criteria`, on PV metadata, configurations, activations, and the v2 `PvQuery.attr` selector. Both rest on server behavior a unit test cannot reach: a unit test asserts the request carries `values == []`, but only a real server distinguishes an existence filter from an `$in: []` that matches nothing. Each test therefore stores an attribute value, asserts the key-only form finds the record, and asserts a query for a *different* value does not -- that pairing is what makes the first assertion meaningful. The v2 class ingests its own samples (the selector needs archived data) through `ingest_support`, and confirms they are queryable with a *name-list* selector before asserting the negative case, so an empty result can only mean the attribute selector matched nothing. Catalogue records are torn down per run, but the ingested samples are not -- the archive has no delete RPC, so each run leaves a few samples under a run-unique PV name, the same residue `test_datasets_annotations_integration.py` leaves @@ -559,8 +560,9 @@ Invariants worth knowing before touching this code (dp-service citations are in quotas; both matched texts live in module constants, since neither is a stable API. - **A frame's column provenance survives ingestion only into a 1.16.0 or later server**; the rest of the ingestion proto is unchanged since rel-1.15.0. -- **Non-scalar columns cannot be read back live yet**: `querySamples` is scalar-only, so array, image, struct, - and serialized data is verified only as far as ingestion (ack + SUCCESS) until the bucket query (#16). +- **Non-scalar columns are read back through the bucket query** (#16), since `querySamples` is scalar-only; + `test_non_scalar_columns_read_back_exactly` checks values, dims, image descriptor, `schemaId`, and the serialized + payload and encoding after ingestion. - **Live verification** (#17 PR B): `tests/integration/test_ingestion_client_integration.py` covers every behavior above that a server can show -- see Key Files -- and `doc/cookbook/ingestion.md` was run as one continuous script against a live stack. That includes the `RESOURCE_EXHAUSTED` hint on both a unary call and a stream, which also @@ -734,6 +736,73 @@ Notes: `column_table_to_torch()` behind a separate optional `[torch]` extra — no change to `QueryClient` or the NumPy path. Not built yet. +### Bucket Query API (Query Service) + +Issue #16 (`plan/tickets/16/plan.md`) wrapped `queryBuckets` / `queryBucketsStream` on the same `client.query`, over +the same `QueryParams`; `bucket_conversions` reads the results. Worked example: the "Whole buckets" section of +`doc/cookbook/query.md`. The server facts below are from dp-service `08c2038` (citations in the plan's T4-T6). + +```python +from datetime import datetime, timezone +from dp_python_lib.client import MldpClient, QueryParams, PvQuery as PV +from dp_python_lib.client import bucket_conversions as bc +from dp_python_lib.client import data_frame_conversions as dfc + +q = MldpClient().query +begin = datetime(2026, 2, 2, 18, 4, 12, 250_000, tzinfo=timezone.utc) +end = datetime(2026, 2, 2, 18, 4, 12, 500_000, tzinfo=timezone.utc) +params = QueryParams(begin_time=begin, end_time=end, pv_selector=PV.name_list(["BPMS:GUNB:314:X"])) + +for page in q.iter_query_buckets(params): # limit counts BUCKETS here + for bucket in page.data_buckets: # whole, untrimmed, in stored form + column = bc.bucket_column(bucket) # any arm, SerializedDataColumn included + values, stamps = bc.bucket_values(bucket), bc.bucket_timestamps(bucket) # integer epoch nanos + dims, meta = dfc.column_dimensions(column), dfc.column_metadata_dict(column) + trimmed = bc.trim_bucket(bucket, begin, end) # exact, half-open; None if nothing is left + +frames = bc.query_buckets_to_dataframes(q, params, trim=True, max_buckets=10_000) # {pv: DataFrame}, [analysis] +``` + +Invariants worth knowing before touching this code: + +- **A bucket comes back whole and verbatim.** Selection is an overlap test (`firstTime < end AND lastTime >= begin`, + half-open), and the stored axis and column are copied unchanged: a `SamplingClock` stays a clock, a typed column + keeps its type and structural fields, a legacy ingest comes back as a `DataColumn`. The column's `name` is the + ingested name, which equals `pvName`. `providerId` / `providerName` are the provider of *that bucket*. +- **Trimming is opt-in and covers `[begin, end)` only.** With `config_criteria`, each activation interval is a + fragment of one `$or` find, so a bucket spanning a gap between intervals is returned once, whole, gap samples + included -- and trimming to the outer range cannot remove them. `query_samples()` is the answer when it matters. +- **Read-side checks only.** Ingestion stores a *non-decreasing* `TimestampList` (repeats allowed; only decreasing + is rejected, since rel-1.13.0), looser than this library's strictly-increasing write rule. So nothing in + `bucket_conversions` calls `validate_data_frame()` or `timestamp_count()`: they would make legitimately stored + buckets unreadable. A *decreasing* axis still raises, because `trim_bucket()`'s `bisect_left` needs order. +- **`limit` counts buckets**, defaults to 10,000 at 0, and is **silently clamped** at 100,000 (both server config). + A page is also cut at the outbound byte budget and then ends early *with* a token, so a short page is not the + last one; a single bucket over the whole budget is an ERROR naming the PV. The page token is position-only and + not bound to the query that made it. The stream never pages: `nextPageToken` is `""` on every message, a token + on the request is rejected, and messages are cut by `limit` and the same budget. An empty result is success. +- **Order is `(pvName, firstTime)` in practice but promised nowhere**, so `buckets_by_pv()` groups and sorts itself + and the live test pins the observed order only so a change is noticed. Overlapping buckets (separate ingests of + one PV) are kept, never deduplicated. +- **`sampleStatusSelector` is rejected by presence** on buckets, so `_build_query_buckets_request()` refuses a + `sample_status_filter` before any RPC; dropping it would return unfiltered data to a caller who thinks it filtered. +- **`useSerializedColumns` is inert on buckets** (contrary to `query.proto` and the dp-grpc cookbook; filed as + osprey-dcs/dp-grpc#167): a column is serialized only if it was *stored* serialized. It is then user data the + bucket query is the only way back to, so it passes through `bucket_column()` and is refused, never dropped, by the + value readers and the pandas conversions. With `time_range`, only a serialized bucket overlapping the range is + refused. +- **`excludeColumnMetadata` is honored on buckets**, unlike samples: this is the only path that reads stored + `ColumnMetadata` (and #17's provenance) back live. A bucket without metadata reports `None` in + `attrs["buckets"]`, never an empty summary, and the per-PV `attrs["column_metadata"]` appears only when every + bucket carries identical metadata. The structural attrs (`enum_ids`, `dimensions`, `image_descriptors`, + `schema_ids`) are kept under `exclude_column_metadata`, since the values cannot be read without them. +- **There is no streaming pandas convenience and no `.to_dataframe()`**: a PV spans pages and stream messages, so + per-message frames would be PV fragments. Collect buckets, then call `buckets_to_dataframes()` once. +- **Live verification** (#16 PR B): `tests/integration/test_query_buckets_integration.py` and the non-scalar + read-back in `test_ingestion_client_integration.py` (see Key Files), and the new and changed recipes in + `query.md` and `ingestion.md` were run as one continuous script against a live stack run from a local dp-service + development checkout, not a release tag. Do not claim verification against a tag the tests did not run against. + ### DataSets, Annotations, and Export API (Annotation Service) Issue #6 (`plan/tickets/6/plan.md`) added three feature clients on the `annotation` facade: @@ -929,8 +998,8 @@ Notes: as a server rejection. Because absence means "no assertion", an *unlabeled* sample never matches: `exclude()` keeps it, `include()` drops it. - `sampleStatusSelector` is supported by `querySamples()`/`querySamplesStream()` **only** — the server rejects it on a - bucket query. `_build_query_spec()` is shared with the future bucket client (#16) and carries a note at the seam: - a bucket request builder must *refuse* `sample_status_filter` rather than copy it through. + bucket query. `_build_query_spec()` is shared with the bucket request builder (#16), which *refuses* a + `sample_status_filter` with a `ValueError` rather than copying it through or dropping it. - `delete_sample_statuses()` is exact at the sample axis (a straddling bucket is not deleted wholesale), and a delete matching nothing is a success with `deleted_count == 0`. - The deferred domain-registry RPCs (`saveSampleStatusDomain` / `querySampleStatusDomains`) are reserved placeholders diff --git a/README.md b/README.md index 8b80427..8002327 100644 --- a/README.md +++ b/README.md @@ -76,6 +76,12 @@ for their service. transparent paging (`query_samples()` / `iter_query_samples()`) and server-streaming (`iter_query_samples_stream()`). Results convert to pandas DataFrames, NumPy arrays, and Excel via the optional `[analysis]` extra. +- **v2 query API (buckets)** — `client.query`. The archive's stored buckets, whole, over the same + `QueryParams`: `query_buckets()`, `iter_query_buckets()`, and `iter_query_buckets_stream()`. Each + bucket keeps its stored column type, time axis, and column metadata, so this is how array, image, + struct, and serialized columns are read back. `bucket_conversions` reads them in plain Python, + trims them exactly to a range on request, and, with the `[analysis]` extra, assembles one pandas + DataFrame per PV. See [whole buckets](doc/cookbook/query.md#whole-buckets-arrays-images-and-stored-metadata). - **DataSets** — `client.annotation.datasets`. Name a region of the archive (time ranges plus the PVs covered over them) so it can be found, annotated, and exported later: `save_dataset()`, `get_dataset()`, `query_datasets()`, `iter_datasets()`, `delete_dataset()`, and a @@ -113,8 +119,6 @@ older than your `dp_python_lib` will not implement everything listed here. The - **Ingestion Service** - `subscribeData()` — receive data for specified PVs from the ingestion stream - **Query Service** - - `queryBuckets()` / `queryBucketsStream()` — raw data buckets - ([issue #16](https://github.com/osprey-dcs/dp-python-lib/issues/16)) - `queryData()` — bucketed PV time-series data - `queryTable()` — PV time-series data in tabular format - `queryPvStats()` — archive ingestion statistics for PVs @@ -202,7 +206,7 @@ the recipes share one continuous worked example drawn from an accelerator facili | [Cataloguing PVs](doc/cookbook/pv-metadata.md) | Recording what a PV is, then finding PVs by property instead of by name | | [Recording machine configuration](doc/cookbook/machine-configuration.md) | Defining configurations, recording when each was active, and answering "what was the machine doing at 18:04?" | | [Ingesting data](doc/cookbook/ingestion.md) | Registering a provider, sending frames of samples, confirming they landed, and chunking data too big for one message | -| [Querying time-series data](doc/cookbook/query.md) | Retrieving samples by PV, metadata, or machine configuration, and converting to pandas / NumPy / Excel | +| [Querying time-series data](doc/cookbook/query.md) | Retrieving samples by PV, metadata, or machine configuration, and converting to pandas / NumPy / Excel; reading whole stored buckets, including array, image, and struct columns | | [Labeling samples](doc/cookbook/sample-status.md) | Recording per-sample status codes, reading them back, and querying data with flagged samples excluded | | [DataSets and annotations](doc/cookbook/datasets-and-annotations.md) | Naming a region of the archive, attaching analysis results with column-level provenance, and exporting | diff --git a/doc/cookbook/README.md b/doc/cookbook/README.md index f13aac0..25d5b7b 100644 --- a/doc/cookbook/README.md +++ b/doc/cookbook/README.md @@ -25,7 +25,7 @@ client. | [Cataloguing PVs](pv-metadata.md) | Recording what a PV *is* — device, area, element type, position — then finding PVs by those properties instead of by name | | [Recording machine configuration](machine-configuration.md) | Defining configurations, recording when each was active, closing and opening intervals, and answering "what was the machine doing at 18:04?" | | [Ingesting data](ingestion.md) | Registering a provider, sending frames of samples, confirming they landed, and chunking and streaming data too big for one message | -| [Querying time-series data](query.md) | Retrieving samples by PV name, by metadata, or by machine configuration, and converting results to pandas / NumPy / Excel | +| [Querying time-series data](query.md) | Retrieving samples by PV name, by metadata, or by machine configuration, and converting results to pandas / NumPy / Excel; reading whole stored buckets, including array, image, and struct columns | | [Labeling samples](sample-status.md) | Recording per-sample status codes, reading them back, and querying data with flagged samples excluded | | [DataSets and annotations](datasets-and-annotations.md) | Naming a region of the archive, attaching analysis results with column-level provenance, round-tripping calculations through pandas, and exporting | diff --git a/doc/cookbook/conventions.md b/doc/cookbook/conventions.md index 6c73812..e58b5f8 100644 --- a/doc/cookbook/conventions.md +++ b/doc/cookbook/conventions.md @@ -296,7 +296,9 @@ continuous. Note that *bucket selection* in data queries is an overlap test rather than containment: a bucket is returned when it overlaps the requested window at all, so boundary buckets may extend past the -range you asked for. Sample-oriented queries (`query_samples()`) trim to the exact range. +range you asked for. Sample-oriented queries (`query_samples()`) trim to the exact range; a +[bucket query](query.md#buckets-come-back-whole) returns those boundary buckets whole, and trimming +them is opt-in. ## Optional dependencies diff --git a/doc/cookbook/ingestion.md b/doc/cookbook/ingestion.md index f48c048..ce21b39 100644 --- a/doc/cookbook/ingestion.md +++ b/doc/cookbook/ingestion.md @@ -28,7 +28,10 @@ from dp_python_lib.client import ( IngestionRequestStatus, RequestStatusQuery as RS, chunked_request_params, + QueryParams, + PvQuery as PV, ) +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 import data_frame_conversions as dfc ``` @@ -323,9 +326,28 @@ frame = dfb.data_frame(axis, [ Array samples may be nested lists or NumPy arrays, one to three dimensions, all the same shape; they are stored flat in row-major order. The image descriptor and the struct's `schema_id` are -per column. These ingest and confirm like any other frame, but **cannot be queried back yet**: -`query_samples()` returns scalar columns only, and the bucket query that returns the rest is -[issue #16](https://github.com/osprey-dcs/dp-python-lib/issues/16). +per column. These ingest and confirm like any other frame. `query_samples()` returns scalar +columns only, so read them back with a bucket query, which returns each column in its stored form: + +```python +# cookbook:partial +params = QueryParams( + begin_time=datetime(2026, 2, 2, 18, 7, tzinfo=timezone.utc), + end_time=datetime(2026, 2, 2, 18, 8, tzinfo=timezone.utc), + pv_selector=PV.name_list(["BPMS:GUNB:314:WAVEFORM", "CAMR:GUNB:100:IMAGE", "BPMS:GUNB:314:STATE"]), +) +for page in client.query.iter_query_buckets(params): + for bucket in page.data_buckets: + print(bucket.pvName, bc.bucket_values(bucket)) + # BPMS:GUNB:314:STATE [b'\x08\x01', b'\x08\x02'] + # BPMS:GUNB:314:WAVEFORM [[0.1, 0.2, 0.3], [0.4, 0.5, 0.6]] + # CAMR:GUNB:100:IMAGE [b'...png bytes...', b'...png bytes...'] +``` + +The dims, image descriptor, and `schema_id` come back on the column; see +[Whole buckets](query.md#whole-buckets-arrays-images-and-stored-metadata). A column ingested with +`dfb.serialized_column()` comes back serialized, payload and encoding intact, and is never decoded +for you. ## Reading failures diff --git a/doc/cookbook/query.md b/doc/cookbook/query.md index 604b0cd..d971f3f 100644 --- a/doc/cookbook/query.md +++ b/doc/cookbook/query.md @@ -1,7 +1,8 @@ # Querying Time-Series Data Retrieving archived PV samples over a time range — by name, by what the PVs *are*, or by what the -machine was *doing* — and getting the results into pandas or NumPy. +machine was *doing* — and getting the results into pandas or NumPy, either as one aligned table of +scalars or as the archive's [whole stored buckets](#whole-buckets-arrays-images-and-stored-metadata). See [API conventions](conventions.md) for result checking and paging. The samples queried here are the ones [Ingesting data](ingestion.md) stores, and the metadata- and configuration-driven @@ -36,6 +37,7 @@ from dp_python_lib.client import query_conversions as qc - [Scoping a query to a machine configuration](#scoping-a-query-to-a-machine-configuration) - [Getting results into pandas and NumPy](#getting-results-into-pandas-and-numpy) - [Large queries: paging and streaming](#large-queries-paging-and-streaming) +- [Whole buckets: arrays, images, and stored metadata](#whole-buckets-arrays-images-and-stored-metadata) - [Also worth knowing](#also-worth-knowing) ## Model @@ -384,13 +386,145 @@ for frame in qc.stream_query_samples_to_dataframes(client.query, params): Remember that **`limit` is the page size**, not a total cap. To stop early, break out of the loop or use `itertools.islice`. +## Whole buckets: arrays, images, and stored metadata + +Everything above is the **sample-oriented** query: one aligned table of scalar columns, trimmed to +the range. The **bucket-oriented** query returns the archive's stored units instead. Each +`DataBucket` holds one PV's column, in the type it was ingested as, over that bucket's own time +axis. Reach for it when you need: + +- **Array, image, struct, or serialized columns.** `query_samples()` returns scalars only, so + this is the way to read back what the [ingestion recipe](ingestion.md#arrays-images-and-structures) + stores. +- **The column metadata stored with the data**, provenance included. The sample path returns + none; a bucket carries its column's `ColumnMetadata`. +- **The stored time axis** — a `SamplingClock` comes back as a clock, not expanded. + +The methods take the same `QueryParams` and mirror the sample methods: +`query_buckets()` (one page), `iter_query_buckets()` (every page), and +`iter_query_buckets_stream()`. The `bucket_conversions` module reads the results; its plain-Python +half needs no extras. + +```python +# cookbook:partial +from dp_python_lib.client import bucket_conversions as bc + +params = QueryParams( + begin_time=datetime(2026, 2, 2, 18, 7, tzinfo=timezone.utc), + end_time=datetime(2026, 2, 2, 18, 8, tzinfo=timezone.utc), + pv_selector=PV.name_list(["BPMS:GUNB:314:WAVEFORM", "CAMR:GUNB:100:IMAGE"]), +) +for page in client.query.iter_query_buckets(params): + for bucket in page.data_buckets: + column = bc.bucket_column(bucket) # the stored column message, whatever its type + print(bucket.pvName, type(column).__name__, bucket.providerName) + print(" at", bc.bucket_timestamps(bucket)) # epoch nanoseconds, exact + print(" values", bc.bucket_values(bucket)) # one entry per sample +``` + +`bucket_values()` gives one entry per sample: a native scalar, an enum's integer code, a flat list +per array sample, or one `bytes` payload per image or struct sample. The fields those samples +cannot be read without stay on the column, through the `data_frame_conversions` accessors: + +```python +# cookbook:partial +for page in client.query.iter_query_buckets(params): + for bucket in page.data_buckets: + column = bc.bucket_column(bucket) + print(dfc.column_dimensions(column)) # [3] for the waveform; None for an image + print(dfc.image_descriptor_dict(column)) # width, height, channels, encoding; None otherwise + print(dfc.column_schema_id(column)) # a struct's schemaId; None otherwise + print(dfc.column_metadata_dict(column)) # tags, attributes, provenance +``` + +`bc.bucket_to_data_frame(bucket)` views a bucket as a one-column `common.DataFrame`, so any of the +[calculations readers](datasets-and-annotations.md) work on it too. + +### Buckets come back whole + +The server returns **every bucket that overlaps `[begin_time, end_time)`, untrimmed**. A query +for a quarter of a second inside a one-second bucket returns the whole second. Trimming is +opt-in, exact, and half-open like `query_samples()`: + +```python +# cookbook:partial +begin = datetime(2026, 2, 2, 18, 4, 12, 250_000, tzinfo=timezone.utc) +end = datetime(2026, 2, 2, 18, 4, 12, 500_000, tzinfo=timezone.utc) +params = QueryParams(begin_time=begin, end_time=end, pv_selector=PV.name_list(["BPMS:GUNB:314:X"])) + +for page in client.query.iter_query_buckets(params): + for bucket in page.data_buckets: + trimmed = bc.trim_bucket(bucket, begin, end) # None if no sample falls in the range + if trimmed is not None: + print(len(bc.bucket_timestamps(bucket)), "->", len(bc.bucket_timestamps(trimmed))) + # 10000 -> 2500 +``` + +A trimmed `SamplingClock` bucket stays a clock, with its start moved to the first sample kept. + +**Trimming cannot remove configuration gaps.** With `config_criteria`, a bucket spanning a gap +between [two activations](#the-result-covers-several-disjoint-intervals) is returned once, whole, +with the samples in the gap still in it, and trimming to `[begin_time, end_time)` leaves them +there. Use `query_samples()` when that matters. + +### One DataFrame per PV + +With the `[analysis]` extra, `query_buckets_to_dataframes()` runs the whole query and returns a +`dict` of PV name to DataFrame. A bucket result is not one table — each PV has its own axis and +possibly its own column type — so there is no single-frame form. + +```python +# cookbook:partial +frames = bc.query_buckets_to_dataframes(client.query, params, trim=True, max_buckets=10_000) +x = frames["BPMS:GUNB:314:X"] +print(len(x), x.attrs["buckets"][0]["provider_name"]) +``` + +Each frame has a UTC index and one column, named for the PV, with the dtype its stored type +implies. `df.attrs["buckets"]` describes every bucket in the frame (time span, sample count, +provider, column metadata); `df.attrs["column_metadata"]` is present only when every bucket agrees. +Array dims, image descriptors, struct schema ids, and enum ids are in `df.attrs` too. + +A PV's buckets are kept as stored: two ingests that overlap in time both appear, so the index can +repeat instants. A PV whose buckets differ in column type or structure — ingested as doubles one +day and int32 the next, say — raises `ValueError` naming both, rather than being coerced. + +### Paging, streaming, and what `limit` counts + +On this path **`limit` counts buckets, not rows**. A page can also end early, with a page token, +once it reaches the server's message-size budget, and the server silently caps `limit` at its +configured maximum (100,000 by default). `iter_query_buckets()` follows the tokens for you. + +The stream has no tokens. A PV's buckets can arrive across several messages, so collect them all +and convert once: + +```python +# cookbook:partial +buckets = [] +for message in client.query.iter_query_buckets_stream(params): + buckets.extend(message.data_buckets) +frames = bc.buckets_to_dataframes(buckets, time_range=(params.begin_timestamp, params.end_timestamp)) +``` + +### Two refusals + +- **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 would hand back data you + believed was filtered. Filter with `query_samples()`, or drop the filter yourself. +- **A serialized column is passed through, never decoded.** It comes back only if it was ingested + as one (`dfb.serialized_column()`). `bc.bucket_column(bucket)` returns it with its `encoding` + and `payload`; `bucket_values()` and the DataFrame conversions raise `ValueError` rather than + drop it. + ## Also worth knowing - **Half-open range.** `[begin_time, end_time)` — a sample exactly at `end_time` is excluded. - Sample-oriented queries trim to the exact range, unlike bucket-oriented ones. -- **Serialized columns are deferred.** The client always requests dense columns + Sample-oriented queries trim to the exact range; bucket-oriented ones return + [whole buckets](#buckets-come-back-whole). +- **Serialized sample results are deferred.** The client always requests dense columns (`useSerializedColumns = False`). A `ColumnTable` carrying `serializedDataColumns` raises - `NotImplementedError` in the conversion layer. + `NotImplementedError` in the conversion layer. (A column *ingested* serialized is a different + thing; read it back with a [bucket query](#two-refusals).) - **Duplicate column names raise.** Both conversions key columns by `DataColumn.name`, so a table with two identically-named columns raises `ValueError` rather than silently dropping one. - **`valueStatus` is gone.** `DataValue.valueStatus` was removed in dp-grpc 1.16.0 (field 15 is @@ -398,8 +532,6 @@ or use `itertools.islice`. flagged samples outright. It was never populated in `querySamples()` results before that. - **An empty result is not an error** — success with an empty table means nothing matched the range and selector. -- **Bucket-oriented queries (`queryBuckets`) are not yet wrapped** by this library; see - [issue #16](https://github.com/osprey-dcs/dp-python-lib/issues/16). ### How far these examples have been verified @@ -414,3 +546,10 @@ real output; and paging, streaming, `query_samples_to_dataframe()`, and `tests/integration/test_query_client_integration.py` pins the data path itself: values and timestamps round-trip exactly, the range is half-open at both bounds, columns on different clocks align with gaps rather than zeros, and pages at a small `limit` concatenate in order. + +The [whole buckets](#whole-buckets-arrays-images-and-stored-metadata) examples were run the same +way, straight after the ingestion recipe's frames: the array, image, and struct read-back printed +what the ingestion recipe shows, and the quarter-second trim cut a 10,000-sample bucket to 2,500. +`tests/integration/test_query_buckets_integration.py` pins the rest: a clock axis returned +verbatim, whole boundary buckets, `limit` counting buckets across pages, the stream matching the +unary pages, and column metadata read back and excluded. diff --git a/doc/cookbook/sample-status.md b/doc/cookbook/sample-status.md index 0978f59..c7aba41 100644 --- a/doc/cookbook/sample-status.md +++ b/doc/cookbook/sample-status.md @@ -368,9 +368,10 @@ delete would remove, run the same range and `(domain, layer)` through `querySampleStatusDomains` are reserved placeholders in the API that return a "not implemented" error, so this client does not wrap them. Domains are for now a convention between producer and consumer, not a registry. -- **Sample status filtering is sample-query only.** `sampleStatusSelector` is rejected on - bucket-oriented queries, which are not yet wrapped by this library - ([issue #16](https://github.com/osprey-dcs/dp-python-lib/issues/16)). +- **Sample status filtering is sample-query only.** The server rejects a status filter on a + [bucket query](query.md#whole-buckets-arrays-images-and-stored-metadata), so `query_buckets()` + and its siblings refuse a `QueryParams` carrying `sample_status_filter` with a `ValueError` + before calling the server, rather than silently returning unfiltered buckets. - **A pandas view of statuses is not built yet.** `sample_status_conversions` returns plain Python objects and needs no optional extras; a DataFrame conversion is a follow-up. diff --git a/doc/release-notes/NEXT.md b/doc/release-notes/NEXT.md index df34872..4717cda 100644 --- a/doc/release-notes/NEXT.md +++ b/doc/release-notes/NEXT.md @@ -190,9 +190,8 @@ dimensions, an `ImageColumn` without a complete image descriptor, a `StructColum anyway, so nothing that used to be accepted end to end is lost. `enum_column()` likewise now rejects a whitespace-only `enum_id`. -Non-scalar columns (arrays, images, structs) can be ingested, but not yet read back through the -query API, which returns scalar columns only; reading them back is -[#16](https://github.com/osprey-dcs/dp-python-lib/issues/16). +Non-scalar columns (arrays, images, structs) are read back with the bucket query, since the sample +query returns scalar columns only; see [Querying whole buckets](#querying-whole-buckets-issue-16). A new cookbook recipe, [Ingesting data](https://github.com/osprey-dcs/dp-python-lib/blob/main/doc/cookbook/ingestion.md), walks through registering, ingesting, confirming, chunking, and streaming. It was run end to end against a live MLDP stack, as @@ -259,6 +258,12 @@ Behaviors worth knowing: structure (say, ingested as double and later as int32) raises, naming both buckets. Per-bucket provider and column metadata is in `df.attrs["buckets"]`. +The [query recipe](https://github.com/osprey-dcs/dp-python-lib/blob/main/doc/cookbook/query.md#whole-buckets-arrays-images-and-stored-metadata) +has a new section on bucket queries, and the [ingestion recipe](https://github.com/osprey-dcs/dp-python-lib/blob/main/doc/cookbook/ingestion.md#arrays-images-and-structures) +now reads its array, image, and struct columns back. Both were run against a live MLDP stack, and +new integration tests read array, image, struct, and serialized columns back exactly, along with +the column metadata stored with a column. + See [#16](https://github.com/osprey-dcs/dp-python-lib/issues/16) and `plan/tickets/16/plan.md`. ## Installing diff --git a/tests/integration/test_ingestion_client_integration.py b/tests/integration/test_ingestion_client_integration.py index 3972657..9823e93 100644 --- a/tests/integration/test_ingestion_client_integration.py +++ b/tests/integration/test_ingestion_client_integration.py @@ -30,7 +30,9 @@ RegisterProviderRequestParams, chunked_request_params, ) +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 import data_frame_conversions as dfc from .ingest_support import INGESTION_ADDRESS, QUERY_ADDRESS, ingest_confirmed, register_provider, require_services @@ -251,24 +253,60 @@ def test_an_oversized_request_on_a_stream_fails_the_call_with_the_split_hint(sel self.assertIsNone(summary.response, "the call failed as a whole, so there is no summary response") # ------------------------------------------------------------------ - # non-scalar columns (verifiable only as far as SUCCESS until #16 wraps queryBuckets) + # non-scalar columns, read back through the bucket query (#16) # ------------------------------------------------------------------ - def test_non_scalar_columns_reach_success(self): + def test_non_scalar_columns_read_back_exactly(self): + # querySamples is scalar-only, so the bucket query is the one way to read these back. Each column is its + # own PV and so its own bucket, returned in its stored form with every structural field intact. + names = {kind: self._pv(kind) for kind in ("waveform", "map", "camera", "struct", "serialized")} frame = dfb.data_frame( dfb.sampling_clock(self.t0, period_nanos=1_000_000, count=2), [ - dfb.double_array_column(self._pv("waveform"), [[1.0, 2.0, 3.0], [4.0, 5.0, 6.0]]), - dfb.int32_array_column(self._pv("map"), [[[1, 2], [3, 4]], [[5, 6], [7, 8]]]), + dfb.double_array_column(names["waveform"], [[1.0, 2.0, 3.0], [4.0, 5.0, 6.0]]), + dfb.int32_array_column(names["map"], [[[1, 2], [3, 4]], [[5, 6], [7, 8]]]), dfb.image_column( - self._pv("camera"), [b"\x00\x01", b"\x02\x03"], width=1, height=2, channels=1, encoding="raw-u8" + names["camera"], [b"\x00\x01", b"\x02\x03"], width=1, height=2, channels=1, encoding="raw-u8" ), - dfb.struct_column(self._pv("struct"), [b"\x0a\x01", b"\x0a\x02"], schema_id="itest:v1"), - dfb.serialized_column(self._pv("serialized"), b"opaque", encoding="itest-bytes"), + dfb.struct_column(names["struct"], [b"\x0a\x01", b"\x0a\x02"], schema_id="itest:v1"), + dfb.serialized_column(names["serialized"], b"opaque", encoding="itest-bytes"), ], ) ingest_confirmed(self.ingestion, self.provider_id, frame) + params = QueryParams( + begin_time=self.t0, + end_time=self.t0 + timedelta(milliseconds=2), + pv_selector=PvQuery.name_list(list(names.values())), + ) + buckets = {b.pvName: b for page in self.client.query.iter_query_buckets(params) for b in page.data_buckets} + self.assertEqual(set(buckets), set(names.values())) + + waveform = bc.bucket_column(buckets[names["waveform"]]) + self.assertEqual(bc.bucket_values(buckets[names["waveform"]]), [[1.0, 2.0, 3.0], [4.0, 5.0, 6.0]]) + self.assertEqual(dfc.column_dimensions(waveform), [3]) + + grid = bc.bucket_column(buckets[names["map"]]) + self.assertEqual(bc.bucket_values(buckets[names["map"]]), [[1, 2, 3, 4], [5, 6, 7, 8]]) # flat per sample + self.assertEqual(dfc.column_dimensions(grid), [2, 2]) + + camera = bc.bucket_column(buckets[names["camera"]]) + self.assertEqual(bc.bucket_values(buckets[names["camera"]]), [b"\x00\x01", b"\x02\x03"]) + self.assertEqual( + dfc.image_descriptor_dict(camera), {"width": 1, "height": 2, "channels": 1, "encoding": "raw-u8"} + ) + + struct = bc.bucket_column(buckets[names["struct"]]) + self.assertEqual(bc.bucket_values(buckets[names["struct"]]), [b"\x0a\x01", b"\x0a\x02"]) + self.assertEqual(dfc.column_schema_id(struct), "itest:v1") + + # Stored serialized, so returned serialized -- whatever useSerializedColumns says (T6). The payload is + # opaque: reachable through bucket_column(), refused by the value readers. + serialized = bc.bucket_column(buckets[names["serialized"]]) + self.assertEqual((serialized.payload, serialized.encoding), (b"opaque", "itest-bytes")) + with self.assertRaisesRegex(ValueError, "SerializedDataColumn"): + bc.bucket_values(buckets[names["serialized"]]) + if __name__ == "__main__": unittest.main(verbosity=2) diff --git a/tests/integration/test_query_buckets_integration.py b/tests/integration/test_query_buckets_integration.py new file mode 100644 index 0000000..2f3f376 --- /dev/null +++ b/tests/integration/test_query_buckets_integration.py @@ -0,0 +1,267 @@ +""" +Live-server tests for the bucket-oriented v2 query (issue #16; plan/tickets/16/plan.md, PR B). + +Every test ingests its own data through ingest_support and waits for the request's SUCCESS status before querying, +so there is no visibility polling: the status document is written after the buckets. One ingest request makes one +bucket per column, which is what lets these tests reason about bucket boundaries exactly. + +What only a live server shows, and so what these pin: + - a bucket comes back whole and in its stored form: a SamplingClock axis verbatim, never trimmed to the query (T4); + - limit counts buckets, and a PV's buckets reassemble exactly across pages (T5); + - the stream returns the same buckets as the unary pages (T5); + - ColumnMetadata and its provenance read back, and excludeColumnMetadata removes it (T6) -- the first live read + of column metadata, since the samples path populates none. + +Prerequisites: a live MLDP ecosystem, ingestion at localhost:50051 and query at localhost:50052; the tests self-skip +without one. Each run ingests under run-unique PV names and leaves those samples behind -- the archive has no +delete RPC -- the same residue the other ingesting integration tests leave. +""" + +import os +import sys +import time +import unittest +import uuid +from datetime import datetime, timezone + +sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) + +from dp_python_lib.client import MldpClient, PvQuery, QueryParams, from_epoch_nanos +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 import data_frame_conversions as dfc + +from .ingest_support import INGESTION_ADDRESS, QUERY_ADDRESS, ingest_confirmed, register_provider, require_services + +try: + import pandas # noqa: F401 -- availability probe for the [analysis] tests + + HAS_PANDAS = True +except ImportError: + HAS_PANDAS = False + + +class TestQueryBucketsIntegration(unittest.TestCase): + PERIOD = 1_000_000 # 1 ms + CLOCK_COUNT = 10 + PAGED_BUCKETS = 3 # per PV, one ingest request apiece + PAGED_COUNT = 5 # samples per paged bucket + + @classmethod + def setUpClass(cls): + require_services(("ingestion", INGESTION_ADDRESS), ("query", QUERY_ADDRESS)) + cls.client = MldpClient() + cls.run_id = uuid.uuid4().hex[:12] + cls.provider_name = f"itest_buckets_{cls.run_id}" + # Whole seconds, so t0's epoch nanos are exact integer arithmetic. + cls.t0 = datetime.fromtimestamp(int(time.time()) - 3600, tz=timezone.utc) + cls.t0_nanos = int(cls.t0.timestamp()) * 1_000_000_000 + + cls.pv_clock = cls._pv("CLOCK") + cls.clock_values = [float(i) for i in range(cls.CLOCK_COUNT)] + cls.pv_a, cls.pv_b = cls._pv("A"), cls._pv("B") + cls.pv_meta = cls._pv("META") + cls.metadata = dfb.column_metadata( + tags=["itest", "derived"], + attributes={"unit": "mm", "run": cls.run_id}, + provenance=dfb.provenance( + source="itest", + process="x2", + derived_from=[dfb.pv_source(cls.pv_clock, (cls.t0, from_epoch_nanos(cls.t0_nanos + 10 * cls.PERIOD)))], + ), + ) + + ingestion = cls.client.ingestion_client + try: + cls.provider_id = register_provider(ingestion, cls.provider_name) + ingest_confirmed(ingestion, cls.provider_id, cls._clock_frame(cls.pv_clock, cls.clock_values)) + # A and B: PAGED_BUCKETS back-to-back buckets each, so a small limit has to page through both PVs. + for pv, offset in ((cls.pv_a, 0.0), (cls.pv_b, 1000.0)): + for n in range(cls.PAGED_BUCKETS): + start = cls.t0_nanos + n * cls.PAGED_COUNT * cls.PERIOD + values = [offset + n * cls.PAGED_COUNT + i for i in range(cls.PAGED_COUNT)] + ingest_confirmed(ingestion, cls.provider_id, cls._clock_frame(pv, values, start)) + meta_frame = dfb.data_frame( + dfb.sampling_clock(cls.t0, period_nanos=cls.PERIOD, count=3), + [dfb.double_column(cls.pv_meta, [0.0, 2.0, 4.0], metadata=cls.metadata)], + ) + ingest_confirmed(ingestion, cls.provider_id, meta_frame) + except (RuntimeError, TimeoutError) as e: + raise unittest.SkipTest(f"could not ingest the bucket test data: {e}") from None + + @classmethod + def _pv(cls, suffix): + return f"ITEST:BUCKETS:{cls.run_id}:{suffix}" + + @classmethod + def _clock_frame(cls, pv, values, start_nanos=None): + start = from_epoch_nanos(cls.t0_nanos if start_nanos is None else start_nanos) + return dfb.data_frame( + dfb.sampling_clock(start, period_nanos=cls.PERIOD, count=len(values)), [dfb.double_column(pv, values)] + ) + + def _params(self, pvs, begin_offset, end_offset, **kwargs): + return QueryParams( + begin_time=from_epoch_nanos(self.t0_nanos + begin_offset), + end_time=from_epoch_nanos(self.t0_nanos + end_offset), + pv_selector=PvQuery.name_list(pvs), + **kwargs, + ) + + def _one_bucket(self, params): + result = self.client.query.query_buckets(params) + self.assertFalse(result.result_status.is_error, result.result_status.message) + self.assertEqual(result.next_page_token, "") + self.assertEqual(len(result.data_buckets), 1) + return result.data_buckets[0] + + @staticmethod + def _key(bucket): + """A bucket's identity for comparing result sets: PV, axis, and values.""" + return bucket.pvName, tuple(bc.bucket_timestamps(bucket)), tuple(bc.bucket_values(bucket)) + + # ------------------------------------------------------------------ + # one bucket: stored form, whole, then trimmed + # ------------------------------------------------------------------ + + def test_a_sampling_clock_bucket_comes_back_verbatim(self): + bucket = self._one_bucket(self._params([self.pv_clock], 0, self.CLOCK_COUNT * self.PERIOD)) + + self.assertEqual(bucket.pvName, self.pv_clock) + self.assertEqual(bucket.providerId, self.provider_id) + self.assertEqual(bucket.providerName, self.provider_name) + self.assertEqual(bucket.dataTimestamps.WhichOneof("value"), "samplingClock") + clock = bucket.dataTimestamps.samplingClock + self.assertEqual((clock.periodNanos, clock.count), (self.PERIOD, self.CLOCK_COUNT)) + self.assertEqual(clock.startTime.epochSeconds * 1_000_000_000 + clock.startTime.nanoseconds, self.t0_nanos) + + column = bc.bucket_column(bucket) + self.assertEqual(type(column).__name__, "DoubleColumn") + self.assertEqual(column.name, self.pv_clock) + self.assertEqual(bc.bucket_values(bucket), self.clock_values) + self.assertEqual( + bc.bucket_timestamps(bucket), [self.t0_nanos + i * self.PERIOD for i in range(self.CLOCK_COUNT)] + ) + + def test_a_sub_window_returns_the_boundary_bucket_whole_and_trim_reduces_it_exactly(self): + # [3 ms, 6 ms) lies inside the one 10-sample bucket: the server returns all 10, untrimmed. + begin, end = 3 * self.PERIOD, 6 * self.PERIOD + bucket = self._one_bucket(self._params([self.pv_clock], begin, end)) + self.assertEqual(bc.bucket_values(bucket), self.clock_values) + + trimmed = bc.trim_bucket(bucket, from_epoch_nanos(self.t0_nanos + begin), from_epoch_nanos(self.t0_nanos + end)) + self.assertIsNotNone(trimmed) + # Half-open at both bounds, and still a clock -- its start shifted to the first retained sample. + self.assertEqual(bc.bucket_values(trimmed), [3.0, 4.0, 5.0]) + self.assertEqual(trimmed.dataTimestamps.WhichOneof("value"), "samplingClock") + self.assertEqual(bc.bucket_timestamps(trimmed), [self.t0_nanos + i * self.PERIOD for i in (3, 4, 5)]) + self.assertEqual(trimmed.providerId, bucket.providerId) + + @unittest.skipUnless(HAS_PANDAS, "requires the [analysis] extra") + def test_query_buckets_to_dataframes_trims_only_when_asked(self): + params = self._params([self.pv_clock], 3 * self.PERIOD, 6 * self.PERIOD) + + whole = bc.query_buckets_to_dataframes(self.client.query, params)[self.pv_clock] + self.assertEqual(whole[self.pv_clock].tolist(), self.clock_values) + + trimmed = bc.query_buckets_to_dataframes(self.client.query, params, trim=True)[self.pv_clock] + self.assertEqual(trimmed[self.pv_clock].tolist(), [3.0, 4.0, 5.0]) + self.assertEqual([ts.value for ts in trimmed.index], [self.t0_nanos + i * self.PERIOD for i in (3, 4, 5)]) + + # ------------------------------------------------------------------ + # paging and streaming + # ------------------------------------------------------------------ + + def _paged_params(self, limit=0): + return self._params([self.pv_a, self.pv_b], 0, self.PAGED_BUCKETS * self.PAGED_COUNT * self.PERIOD, limit=limit) + + def _expected_paged_values(self, offset): + return [offset + i for i in range(self.PAGED_BUCKETS * self.PAGED_COUNT)] + + def test_limit_counts_buckets_and_pvs_reassemble_across_pages(self): + pages = list(self.client.query.iter_query_buckets(self._paged_params(limit=2))) + + # 6 buckets at 2 per page: limit counts buckets, not samples (there are 30 samples). + self.assertEqual([len(page.data_buckets) for page in pages], [2, 2, 2]) + buckets = [bucket for page in pages for bucket in page.data_buckets] + + # The server's observed order is (pvName, firstTime), so each PV's buckets arrive contiguous. The client does + # not rely on that (buckets_by_pv() sorts), but the cookbook describes it, so a change should be noticed. + self.assertEqual([b.pvName for b in buckets], [self.pv_a] * 3 + [self.pv_b] * 3) + + grouped = bc.buckets_by_pv(reversed(buckets)) # deliberately out of order + self.assertEqual(list(grouped), [self.pv_b, self.pv_a]) + for pv, offset in ((self.pv_a, 0.0), (self.pv_b, 1000.0)): + values = [v for bucket in grouped[pv] for v in bc.bucket_values(bucket)] + stamps = [t for bucket in grouped[pv] for t in bc.bucket_timestamps(bucket)] + self.assertEqual(values, self._expected_paged_values(offset)) + self.assertEqual(stamps, [self.t0_nanos + i * self.PERIOD for i in range(len(values))]) + + @unittest.skipUnless(HAS_PANDAS, "requires the [analysis] extra") + def test_paged_buckets_assemble_into_one_dataframe_per_pv(self): + frames = bc.query_buckets_to_dataframes(self.client.query, self._paged_params(limit=2)) + + self.assertEqual(set(frames), {self.pv_a, self.pv_b}) + for pv, offset in ((self.pv_a, 0.0), (self.pv_b, 1000.0)): + df = frames[pv] + self.assertEqual(df[pv].tolist(), self._expected_paged_values(offset)) + self.assertEqual(len(df.attrs["buckets"]), self.PAGED_BUCKETS) + self.assertTrue(all(entry["provider_id"] == self.provider_id for entry in df.attrs["buckets"])) + self.assertTrue(all(entry["column_name"] == pv for entry in df.attrs["buckets"])) + + @unittest.skipUnless(HAS_PANDAS, "requires the [analysis] extra") + def test_max_buckets_stops_paging(self): + with self.assertRaisesRegex(ValueError, "max_buckets=3"): + bc.query_buckets_to_dataframes(self.client.query, self._paged_params(limit=2), max_buckets=3) + + def test_the_stream_returns_the_same_buckets_as_the_unary_pages(self): + unary = [ + b for page in self.client.query.iter_query_buckets(self._paged_params(limit=2)) for b in page.data_buckets + ] + messages = list(self.client.query.iter_query_buckets_stream(self._paged_params(limit=2))) + streamed = [b for message in messages for b in message.data_buckets] + + self.assertTrue(all(message.next_page_token == "" for message in messages)) + self.assertGreater(len(messages), 1) # the stream cuts messages by limit too + self.assertEqual(sorted(map(self._key, streamed)), sorted(map(self._key, unary))) + self.assertEqual(len(streamed), 2 * self.PAGED_BUCKETS) + + def test_an_empty_result_is_a_success_with_no_buckets(self): + result = self.client.query.query_buckets(self._params([self._pv("NEVER-INGESTED")], 0, self.PERIOD)) + self.assertFalse(result.result_status.is_error, result.result_status.message) + self.assertEqual(result.data_buckets, []) + self.assertEqual(result.next_page_token, "") + + # ------------------------------------------------------------------ + # column metadata: the first path that reads it back live + # ------------------------------------------------------------------ + + def test_column_metadata_reads_back_and_is_absent_when_excluded(self): + bucket = self._one_bucket(self._params([self.pv_meta], 0, 3 * self.PERIOD)) + column = bc.bucket_column(bucket) + self.assertTrue(column.HasField("metadata")) + self.assertEqual(column.metadata, self.metadata) + summary = dfc.column_metadata_dict(column) + self.assertEqual(summary["provenance"]["derived_from"][0]["pv_name"], self.pv_clock) + + excluded = self._one_bucket(self._params([self.pv_meta], 0, 3 * self.PERIOD, exclude_column_metadata=True)) + self.assertFalse(bc.bucket_column(excluded).HasField("metadata")) + self.assertEqual(bc.bucket_values(excluded), [0.0, 2.0, 4.0]) + + @unittest.skipUnless(HAS_PANDAS, "requires the [analysis] extra") + def test_dataframe_attrs_carry_metadata_or_none_when_excluded(self): + result = self.client.query.query_buckets(self._params([self.pv_meta], 0, 3 * self.PERIOD)) + df = result.to_dataframes()[self.pv_meta] + self.assertEqual(df.attrs["column_metadata"][self.pv_meta]["tags"], ["itest", "derived"]) + + # The query excluded it, so the buckets carry none -- reported as None, not as an empty summary. + excluded = self.client.query.query_buckets( + self._params([self.pv_meta], 0, 3 * self.PERIOD, exclude_column_metadata=True) + ) + df = excluded.to_dataframes()[self.pv_meta] + self.assertIsNone(df.attrs["buckets"][0]["column_metadata"]) + self.assertNotIn("column_metadata", df.attrs) + + +if __name__ == "__main__": + unittest.main(verbosity=2) From eedae6a14e778c432ffa43eaa107215cd3ad7697 Mon Sep 17 00:00:00 2001 From: Craig McChesney Date: Fri, 2 Oct 2026 10:03:18 -0600 Subject: [PATCH 2/2] fix: import bc and dfc in query.md's imports block (#16 PR B review) The structural-accessors snippet used dfc, which query.md never imported, so the recipe run on its own failed with a NameError. The snippet checker cannot see this because its shared preamble defines every alias. bc moves into the imports block too, from its inline import mid-recipe. Co-Authored-By: Claude Opus 5.5 Claude-Session: https://claude.ai/code/session_01CajTMkjkkzeoWXLSgMpn5k --- doc/cookbook/query.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/doc/cookbook/query.md b/doc/cookbook/query.md index d971f3f..a531b0c 100644 --- a/doc/cookbook/query.md +++ b/doc/cookbook/query.md @@ -27,6 +27,8 @@ from dp_python_lib.client import ( ConfigQuery as CFG, ) from dp_python_lib.client import query_conversions as qc +from dp_python_lib.client import bucket_conversions as bc +from dp_python_lib.client import data_frame_conversions as dfc ``` ## Contents @@ -407,8 +409,6 @@ half needs no extras. ```python # cookbook:partial -from dp_python_lib.client import bucket_conversions as bc - params = QueryParams( begin_time=datetime(2026, 2, 2, 18, 7, tzinfo=timezone.utc), end_time=datetime(2026, 2, 2, 18, 8, tzinfo=timezone.utc),