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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .dev/tools/check-cookbook-snippets.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
79 changes: 74 additions & 5 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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:
Expand Down Expand Up @@ -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
Expand Down
10 changes: 7 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 |

Expand Down
2 changes: 1 addition & 1 deletion doc/cookbook/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 |

Expand Down
4 changes: 3 additions & 1 deletion doc/cookbook/conventions.md
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Loading
Loading