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
23 changes: 22 additions & 1 deletion .dev/tools/check-cookbook-snippets.py
Original file line number Diff line number Diff line change
Expand Up @@ -69,13 +69,19 @@
# without re-establishing it. Adding a name to paper over a broken snippet defeats the point.
PREAMBLE = """\
# --- checker preamble (not part of the recipe) ---
from datetime import datetime, timezone
import contextlib
from datetime import datetime, timedelta, timezone
from typing import Any, Dict, Iterator, List, Optional

from dp_python_lib.client import (
MldpClient,
IngestionClient,
RegisterProviderRequestParams,
IngestDataRequestParams,
IngestionRequestStatus,
RequestStatusQuery,
RequestStatusQuery as RS,
chunked_request_params,
AnnotationClient,
PvMetadataClient,
PvMetadataQuery,
Expand Down Expand Up @@ -162,6 +168,21 @@
annotation_id: str | None = "6aa1bb271a768e97db44d427"
calculations_id: str | None = "6aa1bb271a768e97db44d428"
assert dataset_id is not None and annotation_id is not None and calculations_id is not None

# The ingestion recipe's handles, carried forward the same way: the registration snippet binds provider_id
# (narrowing it with an assert), and every later snippet sends as that provider.
provider_id: str | None = "6aa1bb271a768e97db44d429"
assert provider_id is not None
since: datetime = datetime.now(timezone.utc)
request: IngestDataRequestParams


def acquire(pv: str, count: int) -> list[float]:
# The recipe's stand-in for an acquisition system. Its own snippet is syntax-checked only (no-mypy), since
# this definition and that one would otherwise collide.
return [0.0] * count


# --- end preamble ---
"""

Expand Down
37 changes: 26 additions & 11 deletions CLAUDE.md
Original file line number Diff line number Diff line change
Expand Up @@ -259,8 +259,11 @@ plan documents one change, `CLAUDE.md` documents the invariant it established.
- `tests/unit/test_annotations_client.py` - Unit tests for AnnotationsClient (incl. `calculations()`, absent-vs-empty calculations on save, the `calculations_id` empty-string-vs-None rule, three-tier error handling, paging)
- `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_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) and polls for bucket visibility rather than sleeping; it establishes that visibility 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
- `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
- `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)
Expand Down Expand Up @@ -555,6 +558,10 @@ Invariants worth knowing before touching this code (dp-service citations are in
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).
- **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
pins the server's size-violation text the hint is gated on.

### PV Metadata API (Annotation Service)
PV metadata methods are exposed under the `annotation` facade at `client.annotation.pv_metadata`
Expand Down Expand Up @@ -695,8 +702,11 @@ qc.dataframe_to_excel(df_all, "out.xlsx") # thin to_excel() wrapper (row-

Notes:
- Time inputs accept a tz-aware datetime, epoch seconds, or `common.Timestamp` (shared `to_timestamp()`); `begin`
must be strictly before `end`, and at least one of `pv_selector` / `config_criteria` must be present (a
config-only query is legal). Empty inputs raise `ValueError`, as does a negative `limit` (`limit=0` is
must be strictly before `end`, and **`pv_selector` is required**, with one arm set, even alongside
`config_criteria` (it has no default, so omitting it is a `TypeError`; `None` or an empty `PvSelector()` a `ValueError`): `QuerySpec.pvSelector` is required by the proto, and dp-service rejects its absence
(`querySpec.pvSelector must be specified`, `QueryV2Resolver`). `QueryParams` accepted a config-only query until
#17 PR B, whose cookbook re-verification found the server rejecting it; "every PV under a configuration" is
`PV.pattern(".*")` plus `config_criteria`. Empty inputs raise `ValueError`, as does a negative `limit` (`limit=0` is
meaningful — the server picks a default).
- `PvQuery` (`PV`) selectors: `name_list(values)` / `pattern(str)` / `metadata([...])`, whose criteria are
`pv_name(exact=, prefix=, contains=)` / `aliases(exact=, prefix=, contains=)` (each repeated & coexisting) plus
Expand All @@ -709,7 +719,9 @@ Notes:
`byteArrayValue`→bytes, `imageValue`→`Image(data, file_type)`); an unhandled oneof arm raises. `DataValue.valueStatus`
no longer exists: it was removed in dp-grpc 1.16.0 (field 15 reserved) in favor of the sample status API, and was
never populated in `querySamples()` results before that. Per-column `ColumnMetadata` lands in
`df.attrs["column_metadata"]` unless `exclude_column_metadata=True`.
`df.attrs["column_metadata"]` unless `exclude_column_metadata=True` -- **but the server's sample path populates
none** (dp-service: "the tabular path carries no column metadata"; `excludeColumnMetadata` is inert there), so a
live `querySamples()` result has an empty `df.attrs`. The conversion is exercised only by unit tests.
- Both conversions key columns by `DataColumn.name`, so a `ColumnTable` carrying two columns with the same name
raises `ValueError` rather than silently dropping the earlier one. `column_table_to_numpy()` always returns
1-D arrays: complex arms become 1-D object arrays, so an array-valued column never collapses into a 2-D array
Expand Down Expand Up @@ -783,10 +795,9 @@ Invariants worth knowing before touching this code:
`distinct` on `pvName` over the **buckets** collection (`MongoAnnotationHandler.validateSaveDataSetRequest` →
`MongoSyncQueryClient.executeQueryPvExistence`), so saving PV metadata does **not** satisfy it. Nothing is
validated client-side (the client cannot know what is archived), but any test or example must use an archived PV.
`ingestData()` acks *before* the bucket becomes queryable, so a `saveDataSet` issued immediately after ingesting
still fails; `tests/integration/test_datasets_annotations_integration.py` probes until it succeeds rather than
sleeping a fixed interval, and ingests through the generated stub because it predates the ingestion client;
#17's second PR moves it onto `ingest_data()` and `await_request_statuses()`.
`ingestData()` acks *before* the bucket becomes queryable, so a `saveDataSet` issued immediately after the ack
can still fail; wait for the request's SUCCESS status first, which is written after the buckets.
`tests/integration/test_datasets_annotations_integration.py` does exactly that, through `ingest_support`.
- **The server does not check `begin < end` on a `DataBlock`** (it checks only that each bound is non-zero and that
`pvNames` is non-empty, and never compares the two bounds), so `data_block()`'s check is the only one there is.
- **A `DataBlock`'s range is half-open, `[begin, end)`** — the same convention as the v2 query API's `QueryParams`,
Expand Down Expand Up @@ -821,7 +832,7 @@ Invariants worth knowing before touching this code:
page size (100, hardcoded in `MongoSyncAnnotationClient.DEFAULT_QUERY_LIMIT`, **not** configurable) applies
**unconditionally**: dropping the last criterion does not change the page size, so a bare `query_*()` still
returns one page and a token. That is why `iter_*` is the right call for browsing. Note the v2 query methods
are deliberately *not* included: `QueryParams` still requires a PV selector or config criteria, since a
are deliberately *not* included: `QueryParams` still requires a PV selector, since a
time-series query with no selection is unbounded rather than a browse-all
- `ExportFormat` makes the server-rejected `EXPORT_FORMAT_UNSPECIFIED` unreachable, and `ExportDataRequestParams`
requires at least one of `dataset_id` / `data_blocks` / `calculations_spec`. The exported file lives on the
Expand Down Expand Up @@ -926,8 +937,12 @@ Notes:
Service built from dp-service `main` carrying the 1.16.0 API: exact nanosecond timestamp round-trip through both
axis forms, absent-stays-absent, full-replace upsert, and layer independence. The tests probe for the API first
and skip with an actionable message against a pre-1.16.0 server, since reachability alone does not imply the RPCs
exist. Status *filtering* of query results is still unit-tested only — it needs ingested sample data to attach
to (#17).
exist. Its `TestSampleStatusQueryFiltering` class (#17 PR B) labels ingested samples and asserts the filtered
query: `exclude()` drops exactly the samples with the named code and keeps unlabeled ones, `include()` keeps only
them. Two things re-verifying the cookbook showed: a *dense* label (a model scoring every sample) gives every
sample a status, code 0 included, so an `include()` without `status_codes` keeps them all; and in a multi-PV
query a filtered sample becomes a gap in its own column, the row vanishing only when every column's sample there
is filtered.

### Configuration Priority (High to Low)
1. **Explicit parameters** (direct channels, config objects)
Expand Down
4 changes: 3 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,8 @@ for their service.
that a request passed validation: confirm it landed with `await_request_statuses()` or
`query_request_status()`. `data_frame.split_data_frame()` cuts a large frame into chunks under
the server's message-size and time-span limits, and the `data_frame` builders now cover array,
image, struct, and serialized columns as well as scalars.
image, struct, and serialized columns as well as scalars. See the
[ingestion recipe](doc/cookbook/ingestion.md).

**Supporting framework:** YAML + environment-variable configuration (`MLDP_*`, via
pydantic-settings), TLS-capable channel creation, hierarchical logging, three-tier error handling
Expand Down Expand Up @@ -200,6 +201,7 @@ the recipes share one continuous worked example drawn from an accelerator facili
| [Creating and connecting a client](doc/cookbook/connecting.md) | Building an `MldpClient`, config files and environment variables, TLS, logging |
| [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 |
| [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
13 changes: 6 additions & 7 deletions doc/cookbook/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,16 +24,11 @@ client.
| [Creating and connecting a client](connecting.md) | The four ways to build an `MldpClient`, configuration files and environment variables, TLS, logging, and when sub-clients are `None` |
| [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 |
| [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 |

**Not yet covered: getting data in.** `IngestionClient` currently exposes only
`register_provider()`, so there is no ingestion recipe.
[Issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17) adds the full ingestion client;
when it lands, this cookbook needs an ingestion recipe, and the
[query recipe's examples need re-verifying against real data](query.md#how-far-these-examples-have-been-verified).

## The worked example

The recipes share one continuous example drawn from an accelerator facility, so the data in the
Expand All @@ -44,6 +39,8 @@ query recipes is the data the earlier recipes create:
`Z` and `S`.
- A physics-shift configuration, `cxi-production` (`PATH=CU_HXR`, `E=14.6`, `RATE=10000`,
`MODE=09`), activated over a shift with `DEST=CXI` and `EXP=CXI_3443`.
- Samples of those signals ingested by provider `bpm-daq` during the shift, at 10 kHz from
18:04:12 — the first second of them is what the query and labeling recipes read back.
- Queries that retrieve those PVs by name, by *"every monitor in GUNB"*, and by *"whatever ran
during the CXI shift"*.
- A dataset naming the first hour of that shift, an annotation recording an orbit drift, and a 1 Hz
Expand All @@ -57,7 +54,9 @@ Attribute names and values are the facility's; tag values are illustrative place
the dp-grpc version its stubs were generated from, so a server older than your `dp_python_lib`
will not implement everything documented here. Where a recipe uses something added in a
particular release, it says so in the body — the [v2 query API](query.md) needs a `rel-1.15.0` or
later server, and [sample status](sample-status.md) needs `rel-1.16.0` or later.
later server, [sample status](sample-status.md) needs `rel-1.16.0` or later, and
[column provenance](datasets-and-annotations.md#recording-where-the-numbers-came-from) sent
with [ingested data](ingestion.md#ingesting-a-frame) survives only into `rel-1.16.0` or later.
- Snippets omit imports and client construction except where a recipe is specifically about those
things. Each recipe lists the imports its examples assume.
- Examples check `result_status.is_error` before reading a payload. This is not ceremony: the
Expand Down
Loading
Loading