diff --git a/.dev/tools/check-cookbook-snippets.py b/.dev/tools/check-cookbook-snippets.py index ac74d85..f2e9a13 100755 --- a/.dev/tools/check-cookbook-snippets.py +++ b/.dev/tools/check-cookbook-snippets.py @@ -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, @@ -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 --- """ diff --git a/CLAUDE.md b/CLAUDE.md index a9bf749..ad9e163 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -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) @@ -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` @@ -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 @@ -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 @@ -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`, @@ -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 @@ -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) diff --git a/README.md b/README.md index 9214fed..8b80427 100644 --- a/README.md +++ b/README.md @@ -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 @@ -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 | diff --git a/doc/cookbook/README.md b/doc/cookbook/README.md index 4f79edf..f13aac0 100644 --- a/doc/cookbook/README.md +++ b/doc/cookbook/README.md @@ -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 @@ -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 @@ -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 diff --git a/doc/cookbook/datasets-and-annotations.md b/doc/cookbook/datasets-and-annotations.md index d4adb23..49dce1b 100644 --- a/doc/cookbook/datasets-and-annotations.md +++ b/doc/cookbook/datasets-and-annotations.md @@ -606,11 +606,10 @@ calculations, are left dangling. reuse tokens across queries. - **`patchDataSet` and `patchAnnotation` are not wrapped.** They are reserved placeholders that return "not implemented"; use the full-replace `save_*` methods. -- **Array, image, and struct column builders do not exist yet.** `data_frame.py` covers the typed - scalar columns and the `DataColumn` escape hatch; the rest are - [issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17)'s to design with ingestion data - in hand. Meanwhile, build those column protos directly and pass them to `data_frame()`, which - accepts pre-built columns alongside the ones its builders return. +- **Array, image, struct, and serialized columns have builders too** — `double_array_column()` and + its siblings, `image_column()`, `struct_column()`, and `serialized_column()`; see + [Ingesting data](ingestion.md#arrays-images-and-structures). `data_frame()` still accepts + pre-built column protos alongside the ones the builders return, and checks both the same way. - **Serialized columns are skipped on read.** `data_frame_columns()` ignores `serializedDataColumns`, whose payloads this library does not decode; read the field directly if you need it. @@ -626,9 +625,8 @@ versus id-only on `query_annotations()`, the full-replace clearing behavior, the the referenced-dataset refusal. That test ingests its own samples first, because of the archive-existence rule described in -[Model](#model) — there is no ingestion client yet -([issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17)), so it uses the generated stub -directly. +[Model](#model), and waits for the ingest's SUCCESS status before saving a dataset over them — the +same sequence as [Ingesting data](ingestion.md#confirming-what-landed). The `x_rms_values` in these examples stand in for real analysis output. The **numbers** are illustrative; the calls around them are the verified part. diff --git a/doc/cookbook/ingestion.md b/doc/cookbook/ingestion.md new file mode 100644 index 0000000..f48c048 --- /dev/null +++ b/doc/cookbook/ingestion.md @@ -0,0 +1,369 @@ +# Ingesting Data + +Getting samples into the archive — registering as a data provider, sending frames of time-series +data, and confirming that they actually landed. + +See [API conventions](conventions.md) for result checking and time handling. The frames sent here +are the same `common.DataFrame` an annotation's calculations use, built with the same `data_frame` +builders; [DataSets and annotations](datasets-and-annotations.md#attaching-an-analysis-result) +introduces them. + +All examples use `client.ingestion_client`. + +> The ingestion API has not changed since 1.15.0, so this recipe works against a `rel-1.15.0` +> server, with one exception: a column's metadata and provenance survive ingestion only into a +> `rel-1.16.0` or later one. + +### Imports used by the examples + +```python +# cookbook:skip +import contextlib +from datetime import datetime, timedelta, timezone + +from dp_python_lib.client import ( + MldpClient, + RegisterProviderRequestParams, + IngestDataRequestParams, + IngestionRequestStatus, + RequestStatusQuery as RS, + chunked_request_params, +) +from dp_python_lib.client import data_frame as dfb +from dp_python_lib.client import data_frame_conversions as dfc +``` + +## Contents + +- [Model](#model) — providers, requests, and why an ack is not a receipt +- [Registering a provider](#registering-a-provider) +- [Ingesting a frame](#ingesting-a-frame) +- [Confirming what landed](#confirming-what-landed) +- [Ingesting from pandas](#ingesting-from-pandas) +- [Frames too big for one message](#frames-too-big-for-one-message) — chunking and streaming +- [One result per request: the bidirectional stream](#one-result-per-request-the-bidirectional-stream) +- [Arrays, images, and structures](#arrays-images-and-structures) +- [Reading failures](#reading-failures) +- [Also worth knowing](#also-worth-knowing) + +## Model + +A **provider** is whatever produces the data — a DAQ process, a script, a detector. It registers +once by name and gets back an id that every request carries. + +A **request** is one frame of data: a time axis plus one column per PV, sent with a +`client_request_id` of your choosing (or a generated one). + +The important part is what the service's answer means. **An ack says only that the request passed +validation.** The service checks the request, acks or rejects it, and *then* queues it for +writing. Anything that goes wrong after that — including an unknown provider id — is invisible at +ack time. What tells you the data landed is the request's **status document**, which the service +writes after the data is stored: + +``` +ingest_data() ──> ack or reject "is it well-formed?" + │ + ▼ (queued) + write buckets + │ + ▼ + status document "did it land?" SUCCESS / ERROR / REJECTED +``` + +So every ingest in this recipe ends with a status check, and a SUCCESS status means the data is +stored and queryable. + +## Registering a provider + +```python +# cookbook:partial +registration = client.ingestion_client.register_provider(RegisterProviderRequestParams( + "bpm-daq", + description="GUNB BPM acquisition", + tag_list=["bpm", "gunb"], + attribute_map={"AREA": "GUNB"}, +)) +if registration.result_status.is_error: + raise RuntimeError(registration.result_status.message) + +provider_id = registration.provider_id +assert provider_id is not None # guaranteed once is_error is False; accessors are Optional +print(provider_id, registration.is_new_provider) +``` + +Registration is **idempotent by name**: registering `"bpm-daq"` again returns the same id, with +`is_new_provider` False. So a process can simply register at startup instead of storing its id. + +## Ingesting a frame + +One second of the worked example's BPM, at 10 kHz: a `SamplingClock` axis (start, period, count) +and one column per PV. + +```python +# cookbook:no-mypy +def acquire(pv: str, count: int) -> list[float]: + """Stand-in for your acquisition system: `count` readings of one PV.""" + offset = {"X": 0.0, "Y": 0.5, "TMIT": 1000.0}[pv.rsplit(":", 1)[1]] + return [offset + (i % 100) / 100 for i in range(count)] +``` + +```python +# cookbook:partial +start = datetime(2026, 2, 2, 18, 4, 12, tzinfo=timezone.utc) +count = 10_000 +frame = dfb.data_frame( + dfb.sampling_clock(start, period_nanos=100_000, count=count), + [dfb.double_column(pv, acquire(pv, count)) + for pv in ("BPMS:GUNB:314:X", "BPMS:GUNB:314:Y", "BPMS:GUNB:314:TMIT")], +) + +request = IngestDataRequestParams(provider_id, frame) # client_request_id: a generated uuid4 +since = datetime.now(timezone.utc) # captured BEFORE sending -- see below +ack = client.ingestion_client.ingest_data(request) +if ack.result_status.is_error: + raise RuntimeError(ack.result_status.message) # rejected: nothing was queued +print(ack.client_request_id, ack.num_rows, ack.num_columns) # '' 10000 3 +``` + +`data_frame()` checks the frame's shape as it builds it — every column has one value per +timestamp, names are unique and non-blank — and `IngestDataRequestParams` checks it again, so a +malformed frame raises `ValueError` naming the column rather than being rejected by the server. +Size limits are left to the server; see [frames too big for one message](#frames-too-big-for-one-message). + +Every column builder also takes `metadata=`, a `dfb.column_metadata(...)` carrying tags, +attributes, and provenance, built exactly as +[DataSets and annotations](datasets-and-annotations.md#recording-where-the-numbers-came-from) +shows for calculations. A `rel-1.16.0` or later server stores it with the ingested column; an +older one does not. + +## Confirming what landed + +```python +# cookbook:partial +statuses = client.ingestion_client.await_request_statuses( + provider_id, [request.client_request_id], since=since) + +for document in statuses[request.client_request_id]: + status = IngestionRequestStatus(document.ingestionRequestStatus) + print(status.name, document.statusMessage) # SUCCESS + if status != IngestionRequestStatus.SUCCESS: + raise RuntimeError(f"{request.client_request_id}: {status.name}: {document.statusMessage}") +``` + +`await_request_statuses()` polls until every named request has a status document, then returns +them, keyed by request id. It does not judge them — that is the loop above. A few details: + +- **`since` is required**, and should be captured *before* the first request is sent. The status + query has no paging, so the time floor keeps the response small, and it keeps out documents from + earlier runs that happened to reuse a request id. The floor is backed off a minute for clock + skew between you and the server. +- **Each id maps to a list**, because the server does not enforce unique request ids. With the + default generated ids there is exactly one document per request. +- **A missing document means "unknown", not "success".** A few server paths write none at all — + a request rejected while the service shuts down, say — which is why the wait has a `timeout` + (30 seconds by default) and raises `TimeoutError` naming the ids still missing. +- **Never read a status you did not get from a document.** The status enum's zero value is + SUCCESS, so an unset or defaulted field reads as success. + +To look at statuses directly, `query_request_status()` takes criteria built with +`RequestStatusQuery`, ANDed: + +```python +# cookbook:partial +result = client.ingestion_client.query_request_status([ + RS.provider_id(provider_id), + RS.status([IngestionRequestStatus.ERROR, IngestionRequestStatus.REJECTED]), + RS.time_range(datetime.now(timezone.utc) - timedelta(hours=1)), # the last hour; end omitted: now +]) +for document in result.request_statuses or []: + print(document.requestId, document.statusMessage) +``` + +The time range matches when each request *finished processing*, not the time of its data. One +query cannot ask about several request ids — each `RS.request_id()` criterion holds one, and +criteria AND — which is why `await_request_statuses()` queries by provider and time and matches +the ids itself. + +## Ingesting from pandas + +If the data is already in a pandas DataFrame with a tz-aware `DatetimeIndex`, +`data_frame_from_pandas()` converts it, mapping each column's dtype to a typed column (this needs +the `[analysis]` extra). The next second of the same BPM: + +```python +# cookbook:partial +import pandas as pd + +index = pd.date_range(datetime(2026, 2, 2, 18, 4, 13, tzinfo=timezone.utc), periods=10_000, freq="100us") +table = pd.DataFrame( + {pv: acquire(pv, 10_000) for pv in ("BPMS:GUNB:314:X", "BPMS:GUNB:314:Y", "BPMS:GUNB:314:TMIT")}, + index=index, +) + +request = IngestDataRequestParams(provider_id, dfc.data_frame_from_pandas(table)) +since = datetime.now(timezone.utc) +ack = client.ingestion_client.ingest_data(request) +``` + +Confirm it exactly as above. Two things differ from building the frame yourself: + +- **The time axis is always an explicit `TimestampList`**, one timestamp per row, never an inferred + `SamplingClock` — even for a regular index like this one. That is about 14 more bytes per row on + the wire. For regularly-sampled data you can describe with a clock, the builders are smaller. +- **A `NaN` anywhere is rejected.** A typed column holds exactly one value per timestamp and cannot + express a gap; drop or fill the missing samples first, or put the sparse PV in its own frame. + +## Frames too big for one message + +The service accepts messages up to **4,096,000 bytes** by default — around half a million doubles. +A minute of one 10 kHz PV is 600,000. Over the limit, the call fails with a gRPC +`RESOURCE_EXHAUSTED` error rather than a reject, and on a stream it takes every later request down +with it. + +`split_data_frame()` cuts a frame along its time axis into pieces that each fit, lazily; +`chunked_request_params()` wraps them as requests with correlated ids; and `ingest_data_stream()` +sends them all on one call: + +```python +# cookbook:partial +minute = datetime(2026, 2, 2, 18, 5, tzinfo=timezone.utc) +count = 600_000 # one minute at 10 kHz +frame = dfb.data_frame( + dfb.sampling_clock(minute, period_nanos=100_000, count=count), + [dfb.double_column(pv, acquire(pv, count)) for pv in ("BPMS:GUNB:314:X", "BPMS:GUNB:314:Y")], +) + +chunks = dfb.split_data_frame(frame, max_bytes=dfb.SERVER_DEFAULT_MAX_MESSAGE_BYTES) +since = datetime.now(timezone.utc) +summary = client.ingestion_client.ingest_data_stream( + chunked_request_params(provider_id, chunks, base_request_id="gunb-314-1805")) + +print(summary.num_requests, summary.client_request_ids) +# 3 ['gunb-314-1805-0', 'gunb-314-1805-1', 'gunb-314-1805-2'] +if summary.result_status.is_error: + print("rejected:", summary.rejected_request_ids) + +request_ids = summary.client_request_ids or [] +statuses = client.ingestion_client.await_request_statuses(provider_id, request_ids, since=since) +``` + +Some points about each piece: + +- **`max_bytes` bounds the whole request**, frame and ids included, with room reserved for the + longest ids the library allows. `SERVER_DEFAULT_MAX_MESSAGE_BYTES` is the server's *default*, + exported as a reference; a deployment can configure its own, so pass what yours uses. +- **`split_data_frame()` also takes `max_rows` and `max_span_nanos`.** The service caps one + request's time span (first to last timestamp) at a day by default, so a long capture needs + `max_span_nanos` as well. +- **Chunk *n* is `-`**, so the pieces of one frame are recognizable in their acks and + status documents. Omit `base_request_id` for a generated one. +- **Everything is lazy.** `split_data_frame()` and `chunked_request_params()` are generators, and + `ingest_data_stream()` pulls from them as it sends, so a large frame is never held twice. +- **A reject does not end the stream.** The service checks each request as it arrives and queues + every valid one. If any were rejected, the result is an error that *still carries its + response*: `rejected_request_ids` names the failures, and every other request was accepted. It + does not raise, because raising would hide which ones got through. `num_requests` is then + `None`; the server omits it. +- **If your request generator raises**, that exception is raised from `ingest_data_stream()` as + itself, with a note giving how many requests had already been handed over. Some of those may + already be stored, so check their status before re-sending anything — a re-send of the same data + fails, as [reading failures](#reading-failures) explains. + +## One result per request: the bidirectional stream + +`iter_ingest_data_bidi_stream()` sends requests on one call like `ingest_data_stream()`, but +yields each request's ack or reject as it arrives, in order — useful for a long-running producer +that wants to react to a reject as it happens: + +```python +# cookbook:partial +second = datetime(2026, 2, 2, 18, 6, tzinfo=timezone.utc) + +def one_second_frames(provider_id: str, n: int): + """A producer: n consecutive one-second frames of the BPM's X plane.""" + for i in range(n): + yield IngestDataRequestParams(provider_id, dfb.data_frame( + dfb.sampling_clock(second + timedelta(seconds=i), period_nanos=100_000, count=10_000), + [dfb.double_column("BPMS:GUNB:314:X", acquire("BPMS:GUNB:314:X", 10_000))], + )) + +requests = one_second_frames(provider_id, 10) +with contextlib.closing(client.ingestion_client.iter_ingest_data_bidi_stream(requests)) as results: + for result in results: + if result.result_status.is_error: + print(f"{result.client_request_id} rejected: {result.result_status.message}") + break # leaving the with block cancels the call +``` + +**A reject is yielded, not raised**: it is a fact about one request, and the service carries on +with the rest. Only a transport failure ends the loop, as a `RuntimeError`. + +**To stop early, close the iterator** — `contextlib.closing` does that however the block exits. +Closing cancels the call, so no further requests are sent. A bare `break` without it is not +enough: while anything still refers to the iterator, the library keeps pulling requests from your +producer and sending them, until the iterator is garbage-collected. Reading to the end needs no +close. + +## Arrays, images, and structures + +Beyond scalar columns, the `data_frame` builders cover one fixed-shape array per sample, one +encoded image per sample, one serialized structure per sample, and a whole column as one opaque +payload: + +```python +# cookbook:partial +axis = dfb.sampling_clock(datetime(2026, 2, 2, 18, 7, tzinfo=timezone.utc), period_nanos=1_000_000, count=2) +frame = dfb.data_frame(axis, [ + dfb.double_array_column("BPMS:GUNB:314:WAVEFORM", [[0.1, 0.2, 0.3], [0.4, 0.5, 0.6]]), + dfb.image_column("CAMR:GUNB:100:IMAGE", [b"...png bytes...", b"...png bytes..."], + width=640, height=480, channels=1, encoding="png"), + dfb.struct_column("BPMS:GUNB:314:STATE", [b"\x08\x01", b"\x08\x02"], schema_id="bpm_state:v1"), +]) +``` + +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). + +## Reading failures + +Where a problem shows up depends on who can see it, and when: + +| Where | What | Examples | +|---|---|---| +| `ValueError` when building the params | The library's own checks | a column shorter than the axis, a blank or duplicate name, an id over 256 characters | +| A reject (`is_error` on the ack; `rejected_request_ids` on a stream) | The service's validation | a frame spanning more than a day; a value over a size cap | +| An ERROR status document | Anything after the ack | **an unknown provider id**; a re-ingest of the same data | +| gRPC error, `RESOURCE_EXHAUSTED` | A message over the size limit | a frame that needed `split_data_frame()`; the error message says so | + +Two of these catch people out: + +- **An unknown provider id is acked.** The provider is looked up only when the request is + processed, so a typo in the id succeeds at ack time and fails in the status document. Only the + status check catches it. +- **Re-ingesting the same data fails, possibly halfway.** Stored data is keyed by PV and first + timestamp, so sending a frame again — a retry after a timeout, say — is acked and then ends in + ERROR. The columns are written in order, so a multi-column frame can leave some PVs stored under + an ERROR status. Before retrying, check whether the first attempt landed. + +## Also worth knowing + +- **The client is stricter than the service in two places**: whitespace-only column names and + duplicate timestamps are rejected, although the service accepts both, because the library must + be able to read back what it writes. If you truly need either, build the request by hand and + call the generated stub. +- **Request ids** default to a random uuid4. Pass your own when it should mean something, but keep + it unique: nothing enforces that, and a reused id makes its status ambiguous. +- **There is no delete.** The archive has no RPC for removing ingested samples, which is another + reason to confirm before re-sending. +- **`subscribeData()`**, the service's live subscription to newly ingested data, is not wrapped. + +### How far these examples have been verified + +This recipe was run as one continuous script against a live MLDP stack, with the PV names made +unique per run. The behaviors it describes — the ack-then-ERROR for an unknown provider and for a +re-ingest, the partial reject on a stream, the bidirectional stream's per-request results, +chunked ingest reading back whole, and arrays, images, structs, and serialized columns reaching +SUCCESS — are covered by `tests/integration/test_ingestion_client_integration.py`. diff --git a/doc/cookbook/query.md b/doc/cookbook/query.md index 0efc917..604b0cd 100644 --- a/doc/cookbook/query.md +++ b/doc/cookbook/query.md @@ -3,9 +3,10 @@ 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. -See [API conventions](conventions.md) for result checking and paging. The metadata- and -configuration-driven queries below read the catalogue built in -[Cataloguing PVs](pv-metadata.md) and [Recording machine configuration](machine-configuration.md). +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 +queries read the catalogue built in [Cataloguing PVs](pv-metadata.md) and +[Recording machine configuration](machine-configuration.md). All examples use `client.query`, which is `None` unless a query channel is configured. @@ -51,10 +52,13 @@ PVs are chosen in one of **three mutually exclusive ways** — pick exactly one: | `PV.metadata([...])` | You want PVs by what they *are* — area, type, device | Independently, `config_criteria` restricts results to the intervals when matching machine -configurations were **active**. It can be combined with any PV selector, or used alone — a -config-only query returns everything recorded while that configuration was in effect. +configurations were **active**. It narrows a PV selector; it does not replace one. -At least one of `pv_selector` or `config_criteria` must be present. +**A PV selector is required**, even with `config_criteria`: the server rejects a query without one, +so `QueryParams` refuses one up front. `pv_selector` has no default, so leaving it out raises +`TypeError` (and a type checker flags it); passing `None` or an empty `PvSelector()` raises +`ValueError`. To query every PV, say so with `PV.pattern(".*")` — see +[everything under one configuration](#everything-under-one-configuration). Results arrive as a **`ColumnTable`**: a list of timestamps plus one `DataColumn` per PV. You can work with that directly, or convert it — see @@ -73,7 +77,8 @@ QueryParams(begin_time=end, end_time=begin, pv_selector=PV.name_list(["BPMS:GUNB:314:X"])) # ValueError: begin must precede end ``` -Also rejected: no selector at all, and a negative `limit`. Note that **`limit=0` is meaningful** — +Also rejected: no PV selector (even alongside `config_criteria`), and a negative `limit`. Note +that **`limit=0` is meaningful** — it means "let the server choose a page size". ## Querying a known list of PVs @@ -226,21 +231,22 @@ the outer `[begin_time, end_time)` range. The returned rows are therefore not necessarily contiguous in time. A DataFrame built from them has a jump in its index at each gap, which matters if you resample or difference across it. -### Config-only queries +### Everything under one configuration -Omit the PV selector entirely to get **everything** recorded while a configuration was active: +To get **every PV** recorded while a configuration was active, select all PVs explicitly: ```python # cookbook:partial params = QueryParams( begin_time=datetime(2026, 2, 2, 0, 0, tzinfo=timezone.utc), end_time=datetime(2026, 2, 3, 0, 0, tzinfo=timezone.utc), + pv_selector=PV.pattern(".*"), config_criteria=[CFG.attr("DEST", ["CXI"])], ) ``` -This is legal and occasionally what you want, but it can return a great deal of data. Bound it -with a tight time range and a `limit`. +Occasionally what you want, but it can return a great deal of data: bound it with a tight time +range and a `limit`, and prefer the streaming form below. ## Getting results into pandas and NumPy @@ -265,27 +271,24 @@ df = result.to_dataframe() print(df) ``` -The frame has a **UTC datetime index** and one column per PV, named by `DataColumn.name`: +The frame has a **UTC datetime index** and one column per PV, named by `DataColumn.name`. For the +three BPM signals [the ingestion recipe](ingestion.md#ingesting-a-frame) stores at 10 kHz: ``` - BPMS:GUNB:314:X BPMS:GUNB:314:TMIT -2026-02-02 17:00:00+00:00 0.1 100 -2026-02-02 17:00:01+00:00 0.2 200 -2026-02-02 17:00:02+00:00 0.3 300 + BPMS:GUNB:314:TMIT BPMS:GUNB:314:X BPMS:GUNB:314:Y +2026-02-02 18:04:12+00:00 1000.00 0.00 0.50 +2026-02-02 18:04:12.000100+00:00 1000.01 0.01 0.51 +2026-02-02 18:04:12.000200+00:00 1000.02 0.02 0.52 ``` -Per-column metadata from the catalogue — tags and attributes — travels with the results and lands -in `df.attrs`: +Columns need not come back in the order the selector listed them; here they arrived sorted by name. -```python -# cookbook:partial -df = client.query.query_samples(params).to_dataframe() -print(df.attrs["column_metadata"]) -# {'BPMS:GUNB:314:X': {'tags': ['production'], 'attributes': {'AREA': 'GUNB'}}} -``` - -Pass `exclude_column_metadata=True` to `QueryParams` to skip fetching it, or -`to_dataframe(exclude_column_metadata=True)` to drop it from the frame. +**Sample query results carry no column metadata today.** The conversions put any `ColumnMetadata` +a result carries into `df.attrs["column_metadata"]`, and `QueryParams` has an +`exclude_column_metadata` flag to suppress it, but the server's sample-query path populates none — +not the catalogue's tags and attributes, and not metadata ingested with the columns — so +`df.attrs` is empty. Read the catalogue with +[`get_pv_metadata()`](pv-metadata.md#looking-up-a-single-pv) when you need it alongside the data. ### The whole query to one DataFrame @@ -343,7 +346,9 @@ Three ways to consume a query, in increasing order of scale: **`query_samples()`** — one page. Simple, and enough when you know the result is small. -**`iter_query_samples()`** — pages transparently, yielding one result per page: +**`iter_query_samples()`** — pages transparently, yielding one result per page. The last page can +be empty: a full page comes with a token, and the server discovers there is nothing more only when +asked for the next one. ```python # cookbook:partial @@ -398,18 +403,14 @@ or use `itertools.islice`. ### How far these examples have been verified -Worth knowing what stands behind the recipes on this page, since it differs from the rest of the -cookbook. - -The request-building side was exercised against a live MLDP stack: the queries below are accepted -and well-formed. What has **not** been observed end to end is the data path — rows coming back, -`ColumnTable` populating, a DataFrame with real samples in it. There is currently no way to -ingest sample data from Python (`IngestionClient` exposes only `register_provider()`), so every -query here returned zero rows. The DataFrame and NumPy output shown above is illustrative, and -the conversions themselves are covered by unit tests over hand-built `ColumnTable` objects. - -[Issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17) adds the ingestion client, and -its Phase 3 is a closed-loop ingest→query round-trip. **When that lands, re-verify this recipe -against real data and delete this note** — in particular, confirm the sample DataFrame output -above matches what a real query returns, and add the ingestion recipe this cookbook currently -lacks. +The examples on this page were re-run against a live MLDP stack holding the data +[Ingesting data](ingestion.md) stores: one second of `BPMS:GUNB:314:X`, `:Y`, and `:TMIT` at +10 kHz, catalogued as in [Cataloguing PVs](pv-metadata.md), under a `cxi-production` activation +from 17:00 to 23:00. The name-list, pattern, metadata (`DEVICE`, and `AREA` with `TYPE`), and +configuration (by name and by `EXP`) queries each returned those samples; the DataFrame above is +real output; and paging, streaming, `query_samples_to_dataframe()`, and +`stream_query_samples_to_dataframes()` agreed on the row count. + +`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. diff --git a/doc/cookbook/sample-status.md b/doc/cookbook/sample-status.md index 6657fa5..0978f59 100644 --- a/doc/cookbook/sample-status.md +++ b/doc/cookbook/sample-status.md @@ -86,9 +86,9 @@ those timestamps with `timestamp_list()`. # cookbook:partial # Timestamps taken from the data itself -- never recomputed or rounded. bad_times = [ - datetime(2024, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc), - datetime(2024, 2, 2, 18, 4, 12, 500000, tzinfo=timezone.utc), - datetime(2024, 2, 2, 18, 4, 13, 750000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 500000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 750000, tzinfo=timezone.utc), ] frame = SampleStatusFrame( @@ -107,7 +107,7 @@ frame = SampleStatusFrame( result = client.annotation.sample_status.save_sample_statuses( SaveSampleStatusesRequestParams( frames=[frame], - source="control room, shift log entry 2024-02-02-B", + source="control room, shift log entry 2026-02-02-B", modified_by="operator", ) ) @@ -126,9 +126,9 @@ Several PVs sharing the same bad timestamps go in the same frame, one column eac ```python # cookbook:partial bad_times = [ - datetime(2024, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc), - datetime(2024, 2, 2, 18, 4, 12, 500000, tzinfo=timezone.utc), - datetime(2024, 2, 2, 18, 4, 13, 750000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 500000, tzinfo=timezone.utc), + datetime(2026, 2, 2, 18, 4, 12, 750000, tzinfo=timezone.utc), ] frame = SampleStatusFrame( @@ -161,7 +161,7 @@ def model_confidence() -> list[float]: # and 10,000 confidence values return [1.0] * 10_000 axis = sampling_clock( - start_time=datetime(2024, 2, 2, 18, 4, 12, tzinfo=timezone.utc), + start_time=datetime(2026, 2, 2, 18, 4, 12, tzinfo=timezone.utc), period_nanos=100_000, count=10_000, ) @@ -272,17 +272,28 @@ params = QueryParams( begin_time=begin, end_time=end, pv_selector=PV.name_list(["BPMS:GUNB:314:X"]), - sample_status_filter=SampleStatusFilter.include(domain="ml_anomaly", layers=["ml_model_v1"]), + sample_status_filter=SampleStatusFilter.include( + domain="ml_anomaly", + layers=["ml_model_v1"], + status_codes=[1], # the model's "anomalous" code + ), ) ``` Omitting `layers` matches every layer in the domain; omitting `status_codes` matches any code. +**Name the codes when the labeling is dense.** The model above scored *every* sample, so every +sample carries a status — `0` for the ones it found normal — and an `include()` with no +`status_codes` keeps all of them. The "absence means no assertion" rule decides what happens to unlabeled samples, and it is worth being deliberate about: an unlabeled sample **does not match** the filter. So `exclude()` keeps it (you remove only what was explicitly flagged) and `include()` drops it (you narrow to explicitly labeled samples only). `exclude()` is almost always what you want for analysis. +The filter works per sample, per PV. With several PVs in one query, a sample that is filtered out +becomes a **gap** in its own column — `NaN` in a DataFrame — while the other PVs keep their value +at that instant; a row disappears only when every column's sample there is filtered out. + There is no `MODE_UNSPECIFIED` to fall into — `include()` and `exclude()` are the ways to build a filter, and the server rejects an unspecified mode. If you hand-build a `SampleStatusSelector` instead, `QueryParams` rejects an empty domain or an unspecified mode up front rather than letting @@ -301,7 +312,7 @@ frame = SampleStatusFrame( domain="data_quality", layer="operator_override", data_timestamps=timestamp_list( - [datetime(2024, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc)]), + [datetime(2026, 2, 2, 18, 4, 12, 250000, tzinfo=timezone.utc)]), columns=[SampleStatusColumn( pv_name="BPMS:GUNB:314:X", status_codes=[1], @@ -365,24 +376,15 @@ delete would remove, run the same range and `(domain, layer)` through ### How far these examples have been verified -The save/query/delete loop **has** been exercised against a live Annotation Service built from -dp-service `main` (the 1.16.0 API, pre-release), by -`tests/integration/test_sample_status_client_integration.py`. That covers the parts most likely -to break silently: - -- Timestamps round-trip exactly, through both `timestamp_list()` and `sampling_clock()` — the - dense case labels 100 samples at 1 kHz and checks every expanded position against - `startTime + i * periodNanos`. -- Unsupplied `confidence` / `reasons` come back as `None`, not `0.0` / `""`. -- Re-saving a key with `reasons` omitted clears the stored reason (full replace, not merge). -- The same PV and instant in two layers stay two distinct statuses. - -What has **not** been observed is the last link: labeling *real archived samples* and watching a -status-filtered query drop exactly those rows. The integration tests save statuses and read them -back, but there is no way to ingest sample data from Python yet -([issue #17](https://github.com/osprey-dcs/dp-python-lib/issues/17)), so the statuses they write -have no underlying samples to attach to. - -**When the ingestion client lands, re-verify the filtering examples** — the -[querying with flagged samples removed](#querying-data-with-flagged-samples-removed) section is -the part still standing on unit tests alone. +The examples on this page were re-run against a live MLDP stack, over the BPM samples +[Ingesting data](ingestion.md) stores (one second of `BPMS:GUNB:314:X`, `:Y`, and `:TMIT` at +10 kHz from 18:04:12): the sparse labels and the dense 10,000-sample clock both attach to real +samples, the rows read back with their reasons, `exclude()` drops exactly the samples labeled with +the named code, `include()` keeps only them, and the correction and both deletes behave as +described. + +The tests pin the parts most likely to break silently. +`tests/integration/test_sample_status_client_integration.py` covers exact timestamp round trips +through both axis forms, absent `confidence` / `reasons` coming back as `None`, full-replace +upserts, and layers staying distinct; its `TestSampleStatusQueryFiltering` class labels ingested +samples and asserts which rows a filtered query returns. diff --git a/doc/release-notes/NEXT.md b/doc/release-notes/NEXT.md index 6bd3cdc..a069136 100644 --- a/doc/release-notes/NEXT.md +++ b/doc/release-notes/NEXT.md @@ -193,6 +193,27 @@ Non-scalar columns (arrays, images, structs) can be ingested, but not yet read b query API, which returns scalar columns only; reading them back is [#16](https://github.com/osprey-dcs/dp-python-lib/issues/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 +were the query and sample-status recipes over the data it stores, and the integration tests now +ingest their own data through the library and wait for its SUCCESS status. + +**Upgrade item — an error at call time: `QueryParams` requires a `pv_selector`.** It used to accept +`config_criteria` alone, but the server rejects every such query (`querySpec.pvSelector must be +specified`), so `QueryParams` now refuses it itself. `pv_selector` no longer has a default, so +leaving it out raises `TypeError` (missing argument), which a type checker also reports; passing +`None` or a `PvSelector` with no form set raises `ValueError`. Nothing that worked is lost. To +query every PV under a configuration, select them explicitly with `PvQuery.pattern(".*")` +alongside the `config_criteria`. + +Re-running the older recipes against real data also corrected two things they claimed: + +- **Sample query results carry no column metadata.** The server's sample path populates none, so + `to_dataframe()` leaves `df.attrs` empty; read the catalogue with `get_pv_metadata()` instead. +- **With dense sample-status labeling, name the codes in `SampleStatusFilter.include()`.** A model + that scores every sample gives every sample a status, including its "normal" code, so an + `include()` without `status_codes` keeps them all. + See [#17](https://github.com/osprey-dcs/dp-python-lib/issues/17) and the Ingestion API section of [`CLAUDE.md`](https://github.com/osprey-dcs/dp-python-lib/blob/main/CLAUDE.md). diff --git a/plan/tickets/17/plan.md b/plan/tickets/17/plan.md index 421debb..6efe73d 100644 --- a/plan/tickets/17/plan.md +++ b/plan/tickets/17/plan.md @@ -433,6 +433,17 @@ All dp-service citations are `origin/main` @ `7e8b2e6`, paths relative to - `CLAUDE.md`: drop the remaining "until #17" notes; record which verification now exists. - `doc/release-notes/NEXT.md`: extend #17's section with the verification now done. +*Found in implementation (2026-09-30).* Re-verifying `query.md` and `sample-status.md` against +real data contradicted three of their claims, each confirmed in dp-service and dp-grpc: a +config-only query is rejected (`QuerySpec.pvSelector` is required; `QueryV2Resolver` rejects its +absence), although `QueryParams` accepts one; the sample-query path returns no column metadata, so +`df.attrs["column_metadata"]` is never populated live; and `SampleStatusFilter.include()` without +`status_codes` keeps every densely labeled sample, code 0 included. The docs are corrected here, +and at Craig's direction `QueryParams` now rejects a config-only query (and an empty `PvSelector`) +with a `ValueError` pointing at `PvQuery.pattern(".*")`, although it is a query-client change. The +snippet checker also caught `iter_ingest_data_bidi_stream()` annotated `Iterator`, which made the +documented `contextlib.closing(...)` fail mypy for callers; it is now a `Generator`. + ## Out of scope - **`subscribeData()`** — implemented server-side, but a long-lived bidi subscription with its own diff --git a/src/dp_python_lib/client/ingestion_client.py b/src/dp_python_lib/client/ingestion_client.py index cf58fcc..64b38f2 100644 --- a/src/dp_python_lib/client/ingestion_client.py +++ b/src/dp_python_lib/client/ingestion_client.py @@ -852,7 +852,7 @@ def _send_ingest_data_bidi_stream(self, feed: _RequestFeed) -> Generator[IngestD def iter_ingest_data_bidi_stream( self, requests: Iterable[IngestDataRequestParams] - ) -> Iterator[IngestDataApiResult]: + ) -> Generator[IngestDataApiResult, None, None]: """ Sends ingestion requests on a bidirectional stream, yielding each one's ack or reject as it arrives. @@ -878,7 +878,8 @@ def iter_ingest_data_bidi_stream( Reading to the end needs no close: the call has finished by then. :param requests: The requests to send. - :return: A lazy iterator over per-request results. + :return: A lazy generator over per-request results. Typed as a Generator, not an Iterator, because it has + the close() that contextlib.closing needs. :raises RuntimeError: on a transport or unexpected error. :raises Exception: whatever the iterable raised. """ diff --git a/src/dp_python_lib/client/query_client.py b/src/dp_python_lib/client/query_client.py index 6adb985..d342c95 100644 --- a/src/dp_python_lib/client/query_client.py +++ b/src/dp_python_lib/client/query_client.py @@ -357,16 +357,18 @@ class QueryParams: this release) and, in a future release, the bucket-oriented queryBuckets() requests. A query selects data over a half-open time range [begin_time, end_time) for a set of PVs. The PV set is chosen - by at most one pv_selector form (see PvQuery: name-list, pattern, or metadata) and/or restricted by a list of - config_criteria (see ConfigQuery). A config-only query (all PVs active under a configuration in the window) is - legal, so pv_selector may be omitted -- but at least one of {pv_selector, config_criteria} must be present. + by exactly one pv_selector form (see PvQuery: name-list, pattern, or metadata), optionally restricted by a list + of config_criteria (see ConfigQuery). The selector is required, as the proto says and the server enforces + (dp-service QueryV2Resolver: "querySpec.pvSelector must be specified"): to query every PV active under a + configuration, select them explicitly with PvQuery.pattern(".*"). Through #17 this class accepted a + config-only query, which the server always rejected. """ def __init__( self, begin_time: TimestampInput, end_time: TimestampInput, - pv_selector: query_pb2.PvSelector | None = None, + pv_selector: query_pb2.PvSelector, config_criteria: list["query_pb2.ConfigurationSelector.Criterion"] | None = None, limit: int | None = None, exclude_column_metadata: bool = False, @@ -376,8 +378,8 @@ def __init__( :param begin_time: Inclusive start of the query range (tz-aware datetime, epoch seconds, or common.Timestamp). :param end_time: Exclusive end of the query range (tz-aware datetime, epoch seconds, or common.Timestamp). The range is half-open [begin_time, end_time); the server trims edge samples. - :param pv_selector: The PV selection (see PvQuery: name_list/pattern/metadata). Optional for a config-only - query. At most one form may be set -- the PvQuery helpers each produce exactly one form. + :param pv_selector: The PV selection (see PvQuery: name_list/pattern/metadata). Required, with exactly one + form set -- the PvQuery helpers each produce exactly one form. For every PV, use PvQuery.pattern(".*"). :param config_criteria: List of AND-combined configuration criteria (see ConfigQuery) restricting results to intervals when matching configurations were active. Optional. :param limit: Maximum number of rows to return per page (optional). This is a per-page size, NOT a total @@ -387,8 +389,8 @@ def __init__( :param sample_status_filter: Optional SampleStatusSelector (see SampleStatusFilter.include()/exclude()) restricting results to, or away from, samples carrying matching sample statuses. Supported by querySamples()/querySamplesStream() only -- see the note in _build_query_spec(). - :raises ValueError: if both begin_time and end_time are not supplied, if neither pv_selector nor - config_criteria is present, if begin_time is not strictly before end_time, if limit is negative, or if + :raises ValueError: if both begin_time and end_time are not supplied, if pv_selector is None or sets no + selector form, if begin_time is not strictly before end_time, if limit is negative, or if sample_status_filter is present but carries an empty domain or MODE_UNSPECIFIED. """ if begin_time is None or end_time is None: @@ -402,8 +404,13 @@ def __init__( ): raise ValueError("QueryParams requires begin_time strictly before end_time (half-open [begin, end))") - if pv_selector is None and not config_criteria: - raise ValueError("QueryParams requires at least one of pv_selector or config_criteria") + # The signature makes the selector required, but None or an empty PvSelector() still type-checks at runtime, + # and the server rejects both -- so say so here, with the way to ask for every PV. + if pv_selector is None or pv_selector.WhichOneof("selector") is None: + raise ValueError( + "QueryParams requires a pv_selector (PvQuery.name_list/pattern/metadata); the server rejects a query " + 'without one, even with config_criteria. To select every PV, use PvQuery.pattern(".*")' + ) # Validate eagerly here rather than letting a negative surface as a raw protobuf range error deep in # request building (ExecutionOptions.limit is a uint32). Note limit=0 is explicitly meaningful per the @@ -557,8 +564,7 @@ def _build_query_spec(self, request_params: QueryParams) -> query_pb2.QuerySpec: spec.timeRange.beginTime.CopyFrom(request_params.begin_timestamp) spec.timeRange.endTime.CopyFrom(request_params.end_timestamp) - if request_params.pv_selector is not None: - spec.pvSelector.CopyFrom(request_params.pv_selector) + spec.pvSelector.CopyFrom(request_params.pv_selector) # always set: QueryParams requires it if request_params.config_criteria: spec.configurationSelector.criteria.extend(request_params.config_criteria) diff --git a/tests/integration/ingest_support.py b/tests/integration/ingest_support.py new file mode 100644 index 0000000..9a0adbd --- /dev/null +++ b/tests/integration/ingest_support.py @@ -0,0 +1,66 @@ +""" +Ingesting test data through the library's own IngestionClient, shared by the integration tests. + +Every test that needs archived samples -- a dataset's data block, a v2 query, a sample status to filter on -- +ingests its own, and synchronizes on the request-status document rather than on an ack or a probe loop. An ack +means only that the request passed validation; the status document is written after the buckets, so a SUCCESS +status means the data is persisted and queryable (plan/tickets/17/plan.md, T2). + +It also holds require_services(), the reachability check every ingesting test class runs first. +""" + +import unittest +import uuid +from datetime import datetime, timezone + +import grpc + +from dp_python_lib.client import IngestDataRequestParams, IngestionRequestStatus, RegisterProviderRequestParams + +INGESTION_ADDRESS = "localhost:50051" +QUERY_ADDRESS = "localhost:50052" +ANNOTATION_ADDRESS = "localhost:50053" + + +def require_services(*services: tuple[str, str]) -> None: + """Skips the calling test class unless each (label, address) accepts a connection within 5s.""" + for label, address in services: + channel = grpc.insecure_channel(address) + try: + grpc.channel_ready_future(channel).result(timeout=5) + except grpc.FutureTimeoutError: + raise unittest.SkipTest(f"MLDP {label} service not available at {address}") from None + finally: + channel.close() + + +def register_provider(ingestion, name: str) -> str: + """Registers (or re-finds) a provider by name and returns its id, raising RuntimeError on failure.""" + result = ingestion.register_provider( + RegisterProviderRequestParams(name, description=None, tag_list=None, attribute_map=None) + ) + if result.provider_id is None: + raise RuntimeError(f"could not register ingestion provider {name!r}: {result.result_status.message}") + return result.provider_id + + +def ingest_confirmed(ingestion, provider_id: str, frame, *, timeout: float = 30.0) -> str: + """ + Ingests one frame and waits for its request to reach SUCCESS, returning the client request id. + + :raises RuntimeError: if the request is rejected, or its status document says anything other than SUCCESS. + :raises TimeoutError: if no status document appears within the timeout. + """ + params = IngestDataRequestParams(provider_id, frame, client_request_id=f"itest-{uuid.uuid4()}") + since = datetime.now(timezone.utc) # captured BEFORE sending: the status query's time floor + ack = ingestion.ingest_data(params) + if ack.result_status.is_error: + raise RuntimeError(f"ingest request {params.client_request_id} was rejected: {ack.result_status.message}") + statuses = ingestion.await_request_statuses(provider_id, [params.client_request_id], since=since, timeout=timeout) + for document in statuses[params.client_request_id]: + if document.ingestionRequestStatus != IngestionRequestStatus.SUCCESS: + raise RuntimeError( + f"ingest request {params.client_request_id} ended " + f"{IngestionRequestStatus(document.ingestionRequestStatus).name}: {document.statusMessage}" + ) + return params.client_request_id diff --git a/tests/integration/test_datasets_annotations_integration.py b/tests/integration/test_datasets_annotations_integration.py index 22b1e8c..bc55576 100644 --- a/tests/integration/test_datasets_annotations_integration.py +++ b/tests/integration/test_datasets_annotations_integration.py @@ -24,7 +24,9 @@ ) from dp_python_lib.client.export_client import ExportDataRequestParams, ExportFormat, calculations_spec from dp_python_lib.client.mldp_client import MldpClient -from dp_python_lib.grpc import common_pb2, ingestion_pb2, ingestion_pb2_grpc +from dp_python_lib.grpc import common_pb2 + +from .ingest_support import ingest_confirmed, register_provider def _epoch_nanos(when: datetime) -> int: @@ -50,9 +52,8 @@ class TestDataSetsAnnotationsIntegration(unittest.TestCase): IN THE ARCHIVE. Despite its error text ("no PV metadata found for names: ..."), the server's check is a distinct on pvName over the buckets collection (MongoAnnotationHandler.validateSaveDataSetRequest -> MongoSyncQueryClient.executeQueryPvExistence), so saving PV metadata is not enough -- the PV must have ingested - data. This class therefore ingests a few samples for its run-unique PV in setUpClass, using the generated - ingestion stub directly: the library's IngestionClient wraps only registerProvider() today, and ingestData() is - issue #17. + data. This class therefore ingests a few samples for its run-unique PV in setUpClass, through the library's + IngestionClient, and waits for the request's SUCCESS status before saving anything. Each test writes under a run-unique owner id and tag and deletes what it wrote, so runs neither collide with each other nor with real data. @@ -119,72 +120,21 @@ def _ingest_samples_for_run_pv(cls): Ingests a handful of samples for this run's PV, so data blocks naming it pass saveDataSet's archive-existence check (see the class docstring). - Uses the generated ingestion stub directly rather than the library, whose IngestionClient covers only - registerProvider() today; wrapping ingestData() is issue #17. + Waits for the request's SUCCESS status rather than probing saveDataSet: the status document is written + after the buckets, so once it says SUCCESS the PV is in the archive. """ - channel = grpc.insecure_channel(cls.INGESTION_ADDRESS) - cls._ingestion_channel = channel - stub = ingestion_pb2_grpc.DpIngestionServiceStub(channel) - - registration = stub.registerProvider( - ingestion_pb2.RegisterProviderRequest(providerName=f"itest_provider_{cls.run_id}"), timeout=10 + ingestion = cls.client.ingestion_client + frame = dfb.data_frame( + dfb.sampling_clock(cls.begin_time, period_nanos=cls.SAMPLE_PERIOD_NANOS, count=cls.SAMPLE_COUNT), + [dfb.double_column(cls.pv_name, [float(i) for i in range(cls.SAMPLE_COUNT)])], ) - if registration.HasField("exceptionalResult"): - raise unittest.SkipTest( - f"could not register an ingestion provider: {registration.exceptionalResult.message}" - ) - provider_id = registration.registrationResult.providerId - - request = ingestion_pb2.IngestDataRequest(providerId=provider_id, clientRequestId=f"itest-{cls.run_id}") - clock = request.ingestionDataFrame.dataTimestamps.samplingClock - clock.startTime.epochSeconds = int(cls.begin_time.timestamp()) - clock.periodNanos = cls.SAMPLE_PERIOD_NANOS - clock.count = cls.SAMPLE_COUNT - - column = request.ingestionDataFrame.dataColumns.add() - column.name = cls.pv_name - for i in range(cls.SAMPLE_COUNT): - column.dataValues.add().doubleValue = float(i) - - response = stub.ingestData(request, timeout=15) - if response.HasField("exceptionalResult"): - raise unittest.SkipTest(f"could not ingest test data: {response.exceptionalResult.message}") - - cls._await_pv_in_archive() + try: + provider_id = register_provider(ingestion, f"itest_provider_{cls.run_id}") + ingest_confirmed(ingestion, provider_id, frame) + except (RuntimeError, TimeoutError) as e: + raise unittest.SkipTest(f"could not ingest test data: {e}") from None cls.logger.info("Ingested %d samples for %s", cls.SAMPLE_COUNT, cls.pv_name) - @classmethod - def _await_pv_in_archive(cls, attempts=20, delay_seconds=0.5): - """ - Waits until the ingested PV is visible to saveDataSet's existence check. - - ingestData() acks once the request is accepted, which is before the bucket is committed and queryable, so a - saveDataSet issued immediately afterwards still fails the check. Probe with a throwaway save rather than - sleeping a fixed interval. - """ - probe_params = SaveDataSetRequestParams( - name=f"itest archive probe {cls.run_id}", - owner_id=cls.owner_id, - data_blocks=[data_block(cls.begin_time, cls.end_time, [cls.pv_name])], - ) - for _ in range(attempts): - result = cls.datasets.save_dataset(probe_params) - if not result.result_status.is_error: - cls.datasets.delete_dataset(result.dataset_id) - return - time.sleep(delay_seconds) - - raise unittest.SkipTest( - f"ingested data for {cls.pv_name} did not become visible to saveDataSet within " - f"{attempts * delay_seconds:.0f}s: {result.result_status.message}" - ) - - @classmethod - def tearDownClass(cls): - channel = getattr(cls, "_ingestion_channel", None) - if channel is not None: - channel.close() - @classmethod def _verify_modernized_api_available(cls): """ diff --git a/tests/integration/test_ingestion_client_integration.py b/tests/integration/test_ingestion_client_integration.py index f15038a..3972657 100644 --- a/tests/integration/test_ingestion_client_integration.py +++ b/tests/integration/test_ingestion_client_integration.py @@ -1,171 +1,274 @@ +""" +Live-server tests for IngestionClient (issue #17; plan/tickets/17/plan.md, PR B). + +Every ingest here is closed-loop: an ack says only that a request passed validation, so each test confirms the +outcome through await_request_statuses() -- the request-status document is the one record of whether data landed. +Two tests exist precisely because the ack and the status disagree: an unknown providerId and a re-ingest of the +same PV + first timestamp are both ACKED, and both end in ERROR (T2, T5). + +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 logging import os import sys import time import unittest +import uuid +from datetime import datetime, timedelta, timezone -import grpc - -# Add src directory to path for imports sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) -from dp_python_lib.client.ingestion_client import RegisterProviderRequestParams -from dp_python_lib.client.mldp_client import MldpClient - +from dp_python_lib.client import ( + IngestDataRequestParams, + IngestionRequestStatus, + MldpClient, + PvQuery, + QueryParams, + RegisterProviderRequestParams, + chunked_request_params, +) +from dp_python_lib.client import data_frame as dfb -class TestIngestionClientIntegration(unittest.TestCase): - """ - Integration tests for IngestionClient that require a running MLDP ecosystem. +from .ingest_support import INGESTION_ADDRESS, QUERY_ADDRESS, ingest_confirmed, register_provider, require_services - Prerequisites: - - MLDP services running via docker compose - - Default configuration (localhost:50051 for ingestion service) - - Services should be healthy and accepting connections +# The server's bucket-span cap (Buckets.maxBucketSpanSeconds, 86400 s by default). The client deliberately leaves +# caps to the server, so a frame spanning more than this passes every client check and is rejected by the server +# at validation -- which makes it the reject the stream tests need. +SERVER_MAX_SPAN_SECONDS = 86_400 - To run these tests: - 1. Start MLDP ecosystem: docker compose up -d - 2. Run tests: python -m unittest tests.integration.test_ingestion_client_integration -v - """ +class TestIngestionClientIntegration(unittest.TestCase): @classmethod def setUpClass(cls): - """Set up integration test environment and verify services are reachable.""" - # Configure logging to see detailed operation info logging.basicConfig(level=logging.INFO, format="%(asctime)s - %(name)s - %(levelname)s - %(message)s") - cls.logger = logging.getLogger(__name__) - - cls.logger.info("Setting up integration test environment") - - # Check if services are reachable before running tests - cls._verify_services_available() - - # Create client using default configuration (assumes localhost services) + require_services(("ingestion", INGESTION_ADDRESS), ("query", QUERY_ADDRESS)) cls.client = MldpClient() - cls.logger.info("MldpClient initialized successfully") - - @classmethod - def _verify_services_available(cls): - """Verify that required MLDP services are reachable.""" - cls.logger.info("Checking if MLDP services are available") - - try: - # Try to create a channel to the ingestion service - channel = grpc.insecure_channel("localhost:50051") - - # Set a short timeout for the connection check - grpc.channel_ready_future(channel).result(timeout=5) - cls.logger.info("Ingestion service is reachable at localhost:50051") - channel.close() - - except grpc.FutureTimeoutError: - cls.logger.error("Timeout connecting to ingestion service at localhost:50051") - raise unittest.SkipTest( - "MLDP ingestion service not available at localhost:50051. " - "Please start the MLDP ecosystem with 'docker compose up -d' before running integration tests." - ) from None - except Exception as e: - cls.logger.error("Failed to connect to ingestion service: %s", e) - raise unittest.SkipTest( - f"Cannot connect to MLDP services: {e}. Please ensure the MLDP ecosystem is running." - ) from None - - def test_register_provider_integration_success_or_known_error(self): - """ - Integration test for registerProvider with real gRPC service. - - This test accepts both success and expected business errors (like duplicate provider names) - as valid responses, since the main goal is to verify communication with real services. - """ - self.logger.info("Starting registerProvider integration test") - - # Create unique provider name to avoid conflicts - timestamp = int(time.time()) - unique_name = f"integration-test-provider-{timestamp}" - - # Create test parameters - params = RegisterProviderRequestParams( - name=unique_name, - description="Integration test provider for dp-python-lib", - tag_list=["integration", "test", "automated"], - attribute_map={ - "test_type": "integration", - "framework": "unittest", - "timestamp": str(timestamp), - "client": "dp-python-lib", - }, + cls.ingestion = cls.client.ingestion_client + cls.run_id = uuid.uuid4().hex[:12] + cls.provider_id = register_provider(cls.ingestion, f"itest_ingestion_{cls.run_id}") + # Whole seconds, well in the past, so no test's samples straddle "now" or collide across tests: each test + # takes its own PV name, and a PV's bucket id is its name plus its first timestamp. + cls.t0 = datetime.fromtimestamp(int(time.time()) - 3600, tz=timezone.utc) + + def _pv(self, suffix): + return f"ITEST:INGEST:{self.run_id}:{suffix}" + + def _frame(self, pv, values, start=None, period_nanos=1_000_000): + return dfb.data_frame( + dfb.sampling_clock(start or self.t0, period_nanos=period_nanos, count=len(values)), + [dfb.double_column(pv, list(values))], ) - self.logger.info("Calling registerProvider for provider: %s", unique_name) - - # Call the real service - result = self.client.ingestion_client.register_provider(params) - - # Basic structural assertions - self.assertIsNotNone(result, "Result should not be None") - self.assertIsNotNone(result.result_status, "ResultStatus should not be None") - self.assertIsInstance(result.result_status.is_error, bool, "is_error should be boolean") - self.assertIsInstance(result.result_status.message, str, "message should be string") - - # Log the result for debugging - if result.result_status.is_error: - self.logger.warning("RegisterProvider returned error: %s", result.result_status.message) - - # Check if it's a known/expected error (like duplicate name) - error_msg = result.result_status.message.lower() - known_errors = ["already exists", "duplicate", "conflict", "name is already in use"] - - is_known_error = any(known_err in error_msg for known_err in known_errors) - - if is_known_error: - self.logger.info("Received expected business error - this is normal for integration tests") - else: - # Unexpected error - might indicate real issues - self.logger.error("Unexpected error from registerProvider: %s", result.result_status.message) - # Still don't fail the test - log for investigation - - else: - # Success case - self.logger.info("Successfully registered provider: %s", unique_name) - self.assertIsNotNone(result.response, "Response should not be None on success") - - # The key assertion: we got a proper response structure back - # This proves the gRPC communication, request building, and response parsing all work - self.assertIsNotNone(result, "Integration test passed - received proper response structure") - - def test_register_provider_integration_with_minimal_params(self): - """ - Integration test with minimal required parameters only. - """ - self.logger.info("Starting minimal parameters integration test") - - # Create minimal parameters (only name is required) - timestamp = int(time.time()) - minimal_name = f"minimal-test-{timestamp}" - - params = RegisterProviderRequestParams( - name=minimal_name, - description=None, # Optional - tag_list=None, # Optional - attribute_map=None, # Optional + def _over_span_frame(self, pv): + """Two samples further apart than the server's span cap: valid to the client, rejected by the server.""" + return self._frame(pv, [1.0, 2.0], period_nanos=(SERVER_MAX_SPAN_SECONDS + 1) * 1_000_000_000) + + def _statuses(self, request_ids, since): + """{request id: IngestionRequestStatus}, asserting exactly one document per id.""" + found = self.ingestion.await_request_statuses(self.provider_id, request_ids, since=since) + for request_id, documents in found.items(): + self.assertEqual(len(documents), 1, f"{request_id} has {len(documents)} status documents") + return {rid: IngestionRequestStatus(docs[0].ingestionRequestStatus) for rid, docs in found.items()} + + def _query_values(self, pv, begin, end): + """(epoch nanos, value) pairs for one PV over [begin, end), across every page.""" + params = QueryParams(begin_time=begin, end_time=end, pv_selector=PvQuery.name_list([pv]), limit=1000) + rows = [] + for page in self.client.query.iter_query_samples(params): + table = page.column_table + columns = [c for c in table.dataColumns if c.name == pv] + if not columns: + continue + for ts, value in zip(table.timestampList.timestamps, columns[0].dataValues, strict=True): + rows.append((ts.epochSeconds * 1_000_000_000 + ts.nanoseconds, value.doubleValue)) + return rows + + # ------------------------------------------------------------------ + # registration + # ------------------------------------------------------------------ + + def test_registering_the_same_name_again_returns_the_same_provider(self): + again = self.ingestion.register_provider( + RegisterProviderRequestParams( + f"itest_ingestion_{self.run_id}", description=None, tag_list=None, attribute_map=None + ) + ) + self.assertFalse(again.result_status.is_error, again.result_status.message) + self.assertEqual(again.provider_id, self.provider_id) + self.assertFalse(again.is_new_provider) + + # ------------------------------------------------------------------ + # unary + # ------------------------------------------------------------------ + + def test_unary_ingest_is_acked_and_reaches_success(self): + pv = self._pv("unary") + params = IngestDataRequestParams(self.provider_id, self._frame(pv, [0.1, 0.2, 0.3])) + since = datetime.now(timezone.utc) + ack = self.ingestion.ingest_data(params) + self.assertFalse(ack.result_status.is_error, ack.result_status.message) + self.assertEqual(ack.client_request_id, params.client_request_id) + self.assertEqual((ack.num_rows, ack.num_columns), (3, 1)) + self.assertEqual( + self._statuses([params.client_request_id], since), + {params.client_request_id: IngestionRequestStatus.SUCCESS}, + ) + # SUCCESS means queryable: the values come straight back, nanosecond-exact. + base = int(self.t0.timestamp()) * 1_000_000_000 + self.assertEqual( + self._query_values(pv, self.t0, self.t0 + timedelta(seconds=1)), + [(base, 0.1), (base + 1_000_000, 0.2), (base + 2_000_000, 0.3)], ) - self.logger.info("Calling registerProvider with minimal params for: %s", minimal_name) - - result = self.client.ingestion_client.register_provider(params) + def test_an_unknown_provider_id_is_acked_and_then_ends_in_error(self): + # Contrary to the proto comments (osprey-dcs/dp-grpc#165): the provider is looked up only by the async job. + bogus_provider = uuid.uuid4().hex[:24] + params = IngestDataRequestParams(bogus_provider, self._frame(self._pv("bogus-provider"), [1.0])) + since = datetime.now(timezone.utc) + ack = self.ingestion.ingest_data(params) + self.assertFalse(ack.result_status.is_error, "an unknown providerId is expected to be ACKED") + found = self.ingestion.await_request_statuses(bogus_provider, [params.client_request_id], since=since) + (document,) = found[params.client_request_id] + self.assertEqual(document.ingestionRequestStatus, IngestionRequestStatus.ERROR) + self.assertIn("providerId", document.statusMessage) + + def test_reingesting_the_same_pv_and_first_timestamp_is_acked_and_then_ends_in_error(self): + pv = self._pv("duplicate") + ingest_confirmed(self.ingestion, self.provider_id, self._frame(pv, [1.0, 2.0])) + params = IngestDataRequestParams(self.provider_id, self._frame(pv, [3.0, 4.0])) + since = datetime.now(timezone.utc) + ack = self.ingestion.ingest_data(params) + self.assertFalse(ack.result_status.is_error, "a duplicate bucket is expected to be ACKED") + self.assertEqual( + self._statuses([params.client_request_id], since), + {params.client_request_id: IngestionRequestStatus.ERROR}, + ) + # The first ingest's values stand. + self.assertEqual( + [value for _, value in self._query_values(pv, self.t0, self.t0 + timedelta(seconds=1))], [1.0, 2.0] + ) - # Same structural validation - self.assertIsNotNone(result) - self.assertIsNotNone(result.result_status) + # ------------------------------------------------------------------ + # streaming + # ------------------------------------------------------------------ + + def _three_requests_middle_rejected(self, label): + return [ + IngestDataRequestParams(self.provider_id, self._frame(self._pv(f"{label}-a"), [1.0, 2.0])), + IngestDataRequestParams(self.provider_id, self._over_span_frame(self._pv(f"{label}-b"))), + IngestDataRequestParams(self.provider_id, self._frame(self._pv(f"{label}-c"), [5.0, 6.0])), + ] + + def test_stream_reports_the_rejected_request_and_ingests_the_others(self): + requests = self._three_requests_middle_rejected("stream") + ids = [r.client_request_id for r in requests] + since = datetime.now(timezone.utc) + summary = self.ingestion.ingest_data_stream(requests) + + # A partial reject is an error result that KEEPS its response, and does not raise. + self.assertTrue(summary.result_status.is_error) + self.assertEqual(summary.rejected_request_ids, [ids[1]]) + self.assertEqual(summary.client_request_ids, ids) + self.assertIsNone(summary.num_requests) + self.assertEqual( + self._statuses(ids, since), + { + ids[0]: IngestionRequestStatus.SUCCESS, + ids[1]: IngestionRequestStatus.REJECTED, + ids[2]: IngestionRequestStatus.SUCCESS, + }, + ) - if result.result_status.is_error: - self.logger.info("Minimal params test got error: %s", result.result_status.message) - else: - self.logger.info("Minimal params test succeeded for: %s", minimal_name) + def test_bidi_yields_one_result_per_request_in_order_including_the_reject(self): + requests = self._three_requests_middle_rejected("bidi") + ids = [r.client_request_id for r in requests] + since = datetime.now(timezone.utc) + results = list(self.ingestion.iter_ingest_data_bidi_stream(requests)) + + self.assertEqual([r.client_request_id for r in results], ids) + self.assertEqual([r.result_status.is_error for r in results], [False, True, False]) + self.assertTrue(results[1].result_status.message) + self.assertEqual( + self._statuses(ids, since), + { + ids[0]: IngestionRequestStatus.SUCCESS, + ids[1]: IngestionRequestStatus.REJECTED, + ids[2]: IngestionRequestStatus.SUCCESS, + }, + ) - # Success - we communicated with the service - self.logger.info("Minimal parameters integration test completed") + def test_a_chunked_frame_ingests_as_n_requests_and_reads_back_whole(self): + pv = self._pv("chunked") + n = 2_000 + values = [i * 0.5 for i in range(n)] + frame = self._frame(pv, values) + # Small enough to force several chunks; still room for the two budgeted worst-case ids (2 x 1,024 bytes). + chunks = list(dfb.split_data_frame(frame, max_bytes=8_000)) + self.assertGreater(len(chunks), 2) + + base = f"itest-chunked-{self.run_id}" + since = datetime.now(timezone.utc) + summary = self.ingestion.ingest_data_stream(chunked_request_params(self.provider_id, chunks, base)) + self.assertFalse(summary.result_status.is_error, summary.result_status.message) + self.assertEqual(summary.num_requests, len(chunks)) + ids = [f"{base}-{i}" for i in range(len(chunks))] + self.assertEqual(summary.client_request_ids, ids) + self.assertEqual(set(self._statuses(ids, since).values()), {IngestionRequestStatus.SUCCESS}) + + rows = self._query_values(pv, self.t0, self.t0 + timedelta(seconds=10)) + base_nanos = int(self.t0.timestamp()) * 1_000_000_000 + self.assertEqual(rows, [(base_nanos + i * 1_000_000, v) for i, v in enumerate(values)]) + + # ------------------------------------------------------------------ + # the server's inbound message limit + # ------------------------------------------------------------------ + + def _oversized_frame(self, pv): + """Comfortably over the server's default 4,096,000-byte inbound limit: 600k doubles is ~5.4 MB.""" + return self._frame(pv, [0.0] * 600_000, period_nanos=100_000) + + def test_an_oversized_unary_request_fails_with_the_split_hint(self): + # Over the limit the call fails as RESOURCE_EXHAUSTED -- a transport error, not a reject -- and the hint is + # gated on the server's size-violation text, so this pins that text as well as the hint. + result = self.ingestion.ingest_data( + IngestDataRequestParams(self.provider_id, self._oversized_frame(self._pv("big"))) + ) + self.assertTrue(result.result_status.is_error) + self.assertIn("gRPC error:", result.result_status.message) + self.assertIn("split_data_frame", result.result_status.message) + + def test_an_oversized_request_on_a_stream_fails_the_call_with_the_split_hint(self): + requests = [ + IngestDataRequestParams(self.provider_id, self._frame(self._pv("before-big"), [1.0])), + IngestDataRequestParams(self.provider_id, self._oversized_frame(self._pv("big-streamed"))), + ] + summary = self.ingestion.ingest_data_stream(requests) + self.assertTrue(summary.result_status.is_error) + self.assertIn("split_data_frame", summary.result_status.message) + 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) + # ------------------------------------------------------------------ + + def test_non_scalar_columns_reach_success(self): + 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.image_column( + self._pv("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"), + ], + ) + ingest_confirmed(self.ingestion, self.provider_id, frame) if __name__ == "__main__": - # Allow running this test file directly unittest.main(verbosity=2) diff --git a/tests/integration/test_query_client_integration.py b/tests/integration/test_query_client_integration.py index 2873508..4afd89a 100644 --- a/tests/integration/test_query_client_integration.py +++ b/tests/integration/test_query_client_integration.py @@ -3,14 +3,20 @@ import sys import time import unittest +import uuid +from datetime import datetime, timezone import grpc # Add src directory to path for imports sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) +from dp_python_lib.client import data_frame as dfb from dp_python_lib.client.mldp_client import MldpClient from dp_python_lib.client.query_client import PvQuery, QueryParams +from dp_python_lib.client.time_conversions import from_epoch_nanos + +from .ingest_support import INGESTION_ADDRESS, QUERY_ADDRESS, ingest_confirmed, register_provider, require_services class TestQueryClientIntegration(unittest.TestCase): @@ -18,23 +24,16 @@ class TestQueryClientIntegration(unittest.TestCase): Integration tests for QueryClient that require a running MLDP query service. Prerequisites: - - MLDP query service running (default at localhost:50052), with the v2 query handling enabled. + - MLDP query service running (default at localhost:50052), with the v2 query handling enabled. The closed-loop + class also needs ingestion at localhost:50051. To run these tests: 1. Start the MLDP ecosystem (e.g. docker compose up -d). 2. Run: python -m unittest tests.integration.test_query_client_integration -v - Test-data note: - The query service is read-only, and this library does not yet wrap a data-ingestion RPC, so these tests cannot - create their own queryable data. We deliberately do NOT assume the server is pre-populated with useful data - (a fresh test database has no reason to contain any). Therefore: - - The tests below assert the query *mechanics* end-to-end (the live RPC completes, returns a well-formed - result, tolerates an empty result set, and paging/streaming iterate correctly) without requiring any - specific data to exist. - - The full closed-loop assertions (ingest known data -> query it back -> assert exact values, trimmed - half-open range boundaries, dense alignment on real columns, multi-page paging) are deferred until the - ingestion API client lands. See osprey-dcs/dp-python-lib#17; test_closed_loop_round_trip is skipped with - that reference rather than silently omitted. + The mechanics tests below assume no particular data (a fresh database has none), so they assert only that the + live RPCs complete and return well-formed results. TestQueryClosedLoop ingests its own data and asserts the + values exactly. """ QUERY_ADDRESS = "localhost:50052" @@ -132,22 +131,101 @@ def test_query_samples_stream_mechanics(self): self.assertFalse(page.result_status.is_error) self.logger.info("iter_query_samples_stream received %d message(s)", messages) - @unittest.skip( - "Closed-loop ingest->query round-trip is deferred pending the ingestion API client " - "(osprey-dcs/dp-python-lib#17): ingest a known dataset, query it back, and assert exact value " - "round-trip, trimmed half-open [begin, end) boundaries, dense alignment on real columns, and " - "multi-page paging." - ) - def test_closed_loop_round_trip(self): - # Implementation blocked on the ingestion client (#17). Intended shape: - # 1. Register a provider and ingest a small known dataset (known PV, known timestamps/values) - # over a controlled time window. - # 2. querySamples() over [begin, end) for that PV; assert the exact values and timestamps round-trip. - # 3. Assert half-open trimming: a sample exactly at `end` is excluded, one at `begin` is included. - # 4. Assert dense column/timestamp alignment on the real columns. - # 5. Ingest enough rows to force multiple pages at a small `limit`; assert iter_query_samples() - # concatenates them in order. - raise NotImplementedError + +class TestQueryClosedLoop(unittest.TestCase): + """ + Ingest known data, query it back, and assert it exactly (#17 PR B). + + Two PVs on one run-unique prefix: A every 10 ms and B every 20 ms over the same start, so a query over both + has rows where only A has a sample -- which is what dense alignment has to express. Each ingest is confirmed + through its request-status document before anything is queried, so there is no visibility polling: SUCCESS is + written after the buckets. + """ + + PERIOD_A = 10_000_000 + PERIOD_B = 20_000_000 + COUNT_A = 50 + + @classmethod + def setUpClass(cls): + require_services(("ingestion", INGESTION_ADDRESS), ("query", QUERY_ADDRESS)) + cls.client = MldpClient() + run_id = uuid.uuid4().hex[:12] + cls.pv_a = f"ITEST:QUERY:{run_id}:A" + cls.pv_b = f"ITEST:QUERY:{run_id}:B" + # 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.values_a = [float(i) for i in range(cls.COUNT_A)] + cls.values_b = [100.0 + i for i in range(cls.COUNT_A // 2)] + + ingestion = cls.client.ingestion_client + try: + provider_id = register_provider(ingestion, f"itest_query_{run_id}") + for pv, period, values in ((cls.pv_a, cls.PERIOD_A, cls.values_a), (cls.pv_b, cls.PERIOD_B, cls.values_b)): + frame = dfb.data_frame( + dfb.sampling_clock(cls.t0, period_nanos=period, count=len(values)), [dfb.double_column(pv, values)] + ) + ingest_confirmed(ingestion, provider_id, frame) + except (RuntimeError, TimeoutError) as e: + raise unittest.SkipTest(f"could not ingest the closed-loop test data: {e}") from None + + def _at(self, nanos_after_t0): + return from_epoch_nanos(self.t0_nanos + nanos_after_t0) + + def _params(self, begin_offset, end_offset, pvs=None, limit=0): + return QueryParams( + begin_time=self._at(begin_offset), + end_time=self._at(end_offset), + pv_selector=PvQuery.name_list(pvs or [self.pv_a, self.pv_b]), + limit=limit, + ) + + @staticmethod + def _rows(table): + """{column name: [(epoch nanos, value or None for a gap)]}, asserting dense alignment on the way.""" + stamps = [ts.epochSeconds * 1_000_000_000 + ts.nanoseconds for ts in table.timestampList.timestamps] + out = {} + for column in table.dataColumns: + assert len(column.dataValues) == len(stamps), f"{column.name} is not aligned with the timestamps" + out[column.name] = [ + (t, v.doubleValue if v.WhichOneof("value") is not None else None) + for t, v in zip(stamps, column.dataValues, strict=True) + ] + return out + + def test_values_and_timestamps_round_trip_exactly(self): + result = self.client.query.query_samples(self._params(0, self.COUNT_A * self.PERIOD_A, [self.pv_a])) + self.assertFalse(result.result_status.is_error, result.result_status.message) + self.assertEqual( + self._rows(result.column_table)[self.pv_a], + [(self.t0_nanos + i * self.PERIOD_A, v) for i, v in enumerate(self.values_a)], + ) + + def test_range_is_half_open_at_both_bounds(self): + # [10 ms, 40 ms): the sample exactly at begin is kept, the one exactly at end is not, and neither is the + # one before begin -- both bounds fall inside the ingested bucket, so this is per-sample trimming. + result = self.client.query.query_samples(self._params(self.PERIOD_A, 4 * self.PERIOD_A, [self.pv_a])) + self.assertFalse(result.result_status.is_error, result.result_status.message) + self.assertEqual( + [t - self.t0_nanos for t, _ in self._rows(result.column_table)[self.pv_a]], + [self.PERIOD_A, 2 * self.PERIOD_A, 3 * self.PERIOD_A], + ) + + def test_columns_on_different_clocks_are_densely_aligned(self): + result = self.client.query.query_samples(self._params(0, 6 * self.PERIOD_A)) + self.assertFalse(result.result_status.is_error, result.result_status.message) + rows = self._rows(result.column_table) + self.assertEqual([v for _, v in rows[self.pv_a]], self.values_a[:6]) + # B has a sample on every other row; the rows between are gaps (an unset DataValue), not zeros. + self.assertEqual([v for _, v in rows[self.pv_b]], [100.0, None, 101.0, None, 102.0, None]) + + def test_paging_at_a_small_limit_concatenates_in_order(self): + params = self._params(0, self.COUNT_A * self.PERIOD_A, [self.pv_a], limit=7) + pages = list(self.client.query.iter_query_samples(params)) + self.assertGreater(len(pages), 1) + rows = [row for page in pages for row in self._rows(page.column_table).get(self.pv_a, [])] + self.assertEqual(rows, [(self.t0_nanos + i * self.PERIOD_A, v) for i, v in enumerate(self.values_a)]) if __name__ == "__main__": diff --git a/tests/integration/test_query_helper_relaxations_integration.py b/tests/integration/test_query_helper_relaxations_integration.py index a2368fe..362ca74 100644 --- a/tests/integration/test_query_helper_relaxations_integration.py +++ b/tests/integration/test_query_helper_relaxations_integration.py @@ -37,6 +37,7 @@ # Add src directory to path for imports sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) +from dp_python_lib.client import data_frame as dfb from dp_python_lib.client.machine_config_client import ( ConfigurationActivationQuery, ConfigurationQuery, @@ -46,7 +47,8 @@ from dp_python_lib.client.mldp_client import MldpClient from dp_python_lib.client.pv_metadata_client import PvMetadataQuery, SavePvMetadataRequestParams from dp_python_lib.client.query_client import PvQuery, QueryParams -from dp_python_lib.grpc import ingestion_pb2, ingestion_pb2_grpc + +from .ingest_support import ingest_confirmed, register_provider ANNOTATION_ADDRESS = "localhost:50053" INGESTION_ADDRESS = "localhost:50051" @@ -306,8 +308,7 @@ class TestV2SelectorRelaxations(unittest.TestCase): This is the one path where the client-side non-blank-key check is the ONLY one there is: `QueryV2Resolver` does not validate the key, so a blank one would reach Mongo as an existence test on "attributes." and match nothing silently (`plan/tickets/40/plan.md` T5). Reaching the selector at all needs archived samples, so this - class ingests its own -- through the generated stub, since IngestionClient wraps only registerProvider() - until #17, the same approach test_datasets_annotations_integration.py takes. + class ingests its own through IngestionClient, waiting for the request's SUCCESS status. """ SAMPLE_COUNT = 5 @@ -345,40 +346,21 @@ def tearDownClass(cls): # Guarded: an early skip can leave these unset. if getattr(cls, "client", None) is not None: cls.client.annotation.pv_metadata.delete_pv_metadata(cls.pv_name) - if getattr(cls, "_ingestion_channel", None) is not None: - cls._ingestion_channel.close() @classmethod def _ingest_samples(cls): """Ingests a few samples for this run's PV, so the v2 selector has something to select.""" - cls._ingestion_channel = grpc.insecure_channel(INGESTION_ADDRESS) - stub = ingestion_pb2_grpc.DpIngestionServiceStub(cls._ingestion_channel) - - registration = stub.registerProvider( - ingestion_pb2.RegisterProviderRequest(providerName=f"itest_relax_provider_{cls.run_id}"), timeout=10 + ingestion = cls.client.ingestion_client + frame = dfb.data_frame( + dfb.sampling_clock(cls.begin_time, period_nanos=cls.SAMPLE_PERIOD_NANOS, count=cls.SAMPLE_COUNT), + [dfb.double_column(cls.pv_name, [float(i) for i in range(cls.SAMPLE_COUNT)])], ) - if registration.HasField("exceptionalResult"): - raise unittest.SkipTest( - f"could not register an ingestion provider: {registration.exceptionalResult.message}" - ) - - request = ingestion_pb2.IngestDataRequest( - providerId=registration.registrationResult.providerId, - clientRequestId=f"itest-relax-{cls.run_id}", - ) - clock = request.ingestionDataFrame.dataTimestamps.samplingClock - clock.startTime.epochSeconds = int(cls.begin_time.timestamp()) - clock.periodNanos = cls.SAMPLE_PERIOD_NANOS - clock.count = cls.SAMPLE_COUNT - - column = request.ingestionDataFrame.dataColumns.add() - column.name = cls.pv_name - for i in range(cls.SAMPLE_COUNT): - column.dataValues.add().doubleValue = float(i) - - response = stub.ingestData(request, timeout=15) - if response.HasField("exceptionalResult"): - raise unittest.SkipTest(f"could not ingest test data: {response.exceptionalResult.message}") + try: + provider_id = register_provider(ingestion, f"itest_relax_provider_{cls.run_id}") + # SUCCESS is written after the buckets, so the samples are queryable once this returns. + ingest_confirmed(ingestion, provider_id, frame) + except (RuntimeError, TimeoutError) as e: + raise unittest.SkipTest(f"could not ingest test data: {e}") from None cls.logger.info("Ingested %d samples for %s", cls.SAMPLE_COUNT, cls.pv_name) @classmethod @@ -394,15 +376,12 @@ def _catalogue_pv(cls): if result.result_status.is_error: raise unittest.SkipTest(f"could not catalogue the test PV: {result.result_status.message}") - def _query_columns(self, criterion, attempts=20, delay_seconds=0.5, name_list=False): + def _query_columns(self, criterion, name_list=False): """ Runs a v2 query selecting on `criterion`, returning the column names it produced. - ingestData() acks before the bucket is queryable, so poll rather than sleeping a fixed interval -- the same - reason test_datasets_annotations_integration.py probes for archive visibility. - `name_list=True` ignores `criterion` and selects the run's PV by name instead, which is how the caller - establishes bucket visibility without relying on the attribute selector under test. + establishes that the samples are queryable without relying on the attribute selector under test. """ selector = PvQuery.name_list([self.pv_name]) if name_list else PvQuery.metadata([criterion]) params = QueryParams( @@ -410,27 +389,19 @@ def _query_columns(self, criterion, attempts=20, delay_seconds=0.5, name_list=Fa end_time=self.end_time, pv_selector=selector, ) - for _ in range(attempts): - result = self.client.query.query_samples(params) - self.assertFalse( - result.result_status.is_error, - f"querySamples failed: {result.result_status.message}", - ) - names = [column.name for column in result.column_table.dataColumns] - if names: - return names - time.sleep(delay_seconds) - return [] + result = self.client.query.query_samples(params) + self.assertFalse(result.result_status.is_error, f"querySamples failed: {result.result_status.message}") + return [column.name for column in result.column_table.dataColumns] def test_v2_key_only_attribute_selector_returns_samples(self): - # An empty result means two different things -- "the selector matched nothing" and "the bucket is not - # queryable yet" -- and the negative assertion below reads it as the first. So establish visibility up - # front with a selector that does NOT depend on the behavior under test: a plain name list. Once this - # passes, an empty result from an attribute selector can only be the selector. + # An empty result means two different things -- "the selector matched nothing" and "there is no data" -- + # and the negative assertion below reads it as the first. The ingest waited for SUCCESS, so the samples + # should be queryable; confirm it with a selector that does NOT depend on the behavior under test, a plain + # name list. Once this passes, an empty result from an attribute selector can only be the selector. self.assertIn( self.pv_name, self._query_columns(None, name_list=True), - "the ingested samples never became queryable; cannot distinguish an empty selector result from an " + "the ingested samples are not queryable; cannot distinguish an empty selector result from an " "invisible bucket", ) @@ -450,7 +421,7 @@ def test_v2_key_only_attribute_selector_returns_samples(self): # hits above came from key existence rather than from an unfiltered match. Meaningful because the # visibility assertion above already proved an empty result here is the selector's doing. self.assertEqual( - self._query_columns(PvQuery.attr(self.attribute_key, [NON_MATCHING_VALUE]), attempts=1), + self._query_columns(PvQuery.attr(self.attribute_key, [NON_MATCHING_VALUE])), [], "a value-based v2 selector for a value the PV does not have must select nothing", ) diff --git a/tests/integration/test_sample_status_client_integration.py b/tests/integration/test_sample_status_client_integration.py index 833be90..2281128 100644 --- a/tests/integration/test_sample_status_client_integration.py +++ b/tests/integration/test_sample_status_client_integration.py @@ -3,6 +3,7 @@ import sys import time import unittest +import uuid from datetime import datetime, timedelta, timezone import grpc @@ -10,8 +11,10 @@ # Add src directory to path for imports sys.path.insert(0, os.path.join(os.path.dirname(__file__), "../../src")) +from dp_python_lib.client import data_frame as dfb from dp_python_lib.client import sample_status_conversions as ssc from dp_python_lib.client.mldp_client import MldpClient +from dp_python_lib.client.query_client import PvQuery, QueryParams, SampleStatusFilter from dp_python_lib.client.sample_status_client import ( QuerySampleStatusesRequestParams, SampleStatusColumn, @@ -20,6 +23,16 @@ sampling_clock, timestamp_list, ) +from dp_python_lib.client.time_conversions import from_epoch_nanos + +from .ingest_support import ( + ANNOTATION_ADDRESS, + INGESTION_ADDRESS, + QUERY_ADDRESS, + ingest_confirmed, + register_provider, + require_services, +) _EPOCH = datetime(1970, 1, 1, tzinfo=timezone.utc) @@ -363,5 +376,110 @@ def test_query_matching_nothing_is_success(self): self.assertEqual(result.sample_status_buckets, []) +class TestSampleStatusQueryFiltering(unittest.TestCase): + """ + A sample status filter on a live v2 query (#17 PR B): label some ingested samples, then query the data. + + Six samples a second apart. Samples 1 and 4 are labeled code 2 ("bad"), sample 2 code 1 ("suspect"), the + rest not at all. exclude(code 2) must drop exactly samples 1 and 4 -- keeping the unlabeled ones, since an + absent status asserts nothing -- and include(code 2) must keep only them. The filter names the status code, + so sample 2's other code proves the match is by code and not merely by "has a status". + + Needs ingestion (localhost:50051), query (localhost:50052), and annotation (localhost:50053). The statuses are + deleted afterwards; the ingested samples stay, under a run-unique PV name, as with every ingesting test. + """ + + BAD = 2 + SUSPECT = 1 + COUNT = 6 + PERIOD_NANOS = 1_000_000_000 + + @classmethod + def setUpClass(cls): + require_services(("ingestion", INGESTION_ADDRESS), ("query", QUERY_ADDRESS), ("annotation", ANNOTATION_ADDRESS)) + + cls.client = MldpClient() + cls.sample_status = cls.client.annotation.sample_status + run_id = uuid.uuid4().hex[:12] + cls.pv_name = f"ITEST:SAMPLE:FILTER:{run_id}" + cls.domain = f"itest_filter_domain_{run_id}" + cls.layer = f"itest_filter_layer_{run_id}" + cls.t0 = datetime.fromtimestamp(int(time.time()) - 3600, tz=timezone.utc) + cls.t0_nanos = _exact_nanos(cls.t0) + cls.end = cls.t0 + timedelta(seconds=cls.COUNT) + + ingestion = cls.client.ingestion_client + frame = dfb.data_frame( + sampling_clock(cls.t0, period_nanos=cls.PERIOD_NANOS, count=cls.COUNT), + [dfb.double_column(cls.pv_name, [float(i) for i in range(cls.COUNT)])], + ) + try: + ingest_confirmed(ingestion, register_provider(ingestion, f"itest_filter_{run_id}"), frame) + except (RuntimeError, TimeoutError) as e: + raise unittest.SkipTest(f"could not ingest test data: {e}") from None + + labeled = {1: cls.BAD, 2: cls.SUSPECT, 4: cls.BAD} + saved = cls.sample_status.save_sample_statuses( + SaveSampleStatusesRequestParams( + frames=[ + SampleStatusFrame( + domain=cls.domain, + layer=cls.layer, + data_timestamps=timestamp_list( + [from_epoch_nanos(cls.t0_nanos + i * cls.PERIOD_NANOS) for i in labeled] + ), + columns=[SampleStatusColumn(cls.pv_name, status_codes=list(labeled.values()))], + ) + ], + source="dp-python-lib integration test", + modified_by="dp-python-lib-integration-test", + ) + ) + if saved.result_status.is_error: + raise unittest.SkipTest(f"could not save sample statuses: {saved.result_status.message}") + + @classmethod + def tearDownClass(cls): + if getattr(cls, "sample_status", None) is not None: + cls.sample_status.delete_sample_statuses( + begin_time=cls.t0, end_time=cls.end, domain=cls.domain, layer=cls.layer, pv_names=[cls.pv_name] + ) + + def _sample_indexes(self, sample_status_filter=None): + """The indexes (0..COUNT-1) of the samples a query over the whole range returns.""" + params = QueryParams( + begin_time=self.t0, + end_time=self.end, + pv_selector=PvQuery.name_list([self.pv_name]), + sample_status_filter=sample_status_filter, + ) + result = self.client.query.query_samples(params) + self.assertFalse(result.result_status.is_error, result.result_status.message) + table = result.column_table + (column,) = [c for c in table.dataColumns if c.name == self.pv_name] or [None] + if column is None: + return [] + stamps = [ts.epochSeconds * 1_000_000_000 + ts.nanoseconds for ts in table.timestampList.timestamps] + return [ + (t - self.t0_nanos) // self.PERIOD_NANOS + for t, v in zip(stamps, column.dataValues, strict=True) + if v.WhichOneof("value") is not None + ] + + def test_unfiltered_query_returns_every_sample(self): + self.assertEqual(self._sample_indexes(), list(range(self.COUNT))) + + def test_exclude_drops_exactly_the_matching_samples(self): + self.assertEqual( + self._sample_indexes(SampleStatusFilter.exclude(self.domain, status_codes=[self.BAD])), [0, 2, 3, 5] + ) + + def test_include_keeps_only_the_matching_samples(self): + self.assertEqual(self._sample_indexes(SampleStatusFilter.include(self.domain, status_codes=[self.BAD])), [1, 4]) + + def test_include_without_codes_keeps_every_labeled_sample(self): + self.assertEqual(self._sample_indexes(SampleStatusFilter.include(self.domain)), [1, 2, 4]) + + if __name__ == "__main__": unittest.main() diff --git a/tests/unit/test_query_client.py b/tests/unit/test_query_client.py index 89dd1f9..72a5925 100644 --- a/tests/unit/test_query_client.py +++ b/tests/unit/test_query_client.py @@ -190,14 +190,30 @@ def test_valid_with_pv_selector(self): p = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern("ABC:.*")) self.assertIsNotNone(p.pv_selector) - def test_valid_config_only(self): - # A config-only query (no pv_selector) is legal. - p = QueryParams(BEGIN, END, config_criteria=[ConfigQuery.configuration_name(["c"])]) - self.assertIsNone(p.pv_selector) + def test_omitting_the_selector_is_a_type_error(self): + # The realistic mistake: config criteria alone. pv_selector has no default, so this fails at the call, + # before QueryParams' own check runs -- and mypy flags it statically. + with self.assertRaises(TypeError): + QueryParams(BEGIN, END, config_criteria=[ConfigQuery.configuration_name(["c"])]) # type: ignore[call-arg] + + def test_config_only_is_rejected(self): + # The server requires a PV selector even with config criteria (#17 PR B), so fail here, pointing at ".*". + with self.assertRaises(ValueError) as ctx: + QueryParams(BEGIN, END, pv_selector=None, config_criteria=[ConfigQuery.configuration_name(["c"])]) + self.assertIn('PvQuery.pattern(".*")', str(ctx.exception)) + + def test_requires_a_selector(self): + with self.assertRaises(ValueError): + QueryParams(BEGIN, END, pv_selector=None) - def test_requires_a_selector_or_config(self): + def test_an_empty_selector_is_rejected(self): + # A default-constructed PvSelector sets no oneof arm; the server rejects that too. with self.assertRaises(ValueError): - QueryParams(BEGIN, END) + QueryParams(BEGIN, END, pv_selector=query_pb2.PvSelector()) + + def test_config_criteria_still_narrow_a_selector(self): + p = QueryParams(BEGIN, END, pv_selector=PvQuery.pattern(".*"), config_criteria=[ConfigQuery.category(["c"])]) + self.assertEqual(len(p.config_criteria), 1) def test_begin_equal_end_raises(self): with self.assertRaises(ValueError): @@ -268,10 +284,12 @@ def test_build_spec_roundtrip_full(self): self.assertFalse(req.resultRepresentation.useSerializedColumns) self.assertFalse(req.resultRepresentation.excludeColumnMetadata) - def test_build_config_only(self): - p = QueryParams(BEGIN, END, config_criteria=[ConfigQuery.category(["optics"])]) + def test_build_with_config_criteria(self): + p = QueryParams( + BEGIN, END, pv_selector=PvQuery.pattern(".*"), config_criteria=[ConfigQuery.category(["optics"])] + ) req = self.client._build_query_samples_request(p) - self.assertEqual(req.querySpec.pvSelector.WhichOneof("selector"), None) + self.assertEqual(req.querySpec.pvSelector.pvNamePattern.pattern, ".*") self.assertEqual(len(req.querySpec.configurationSelector.criteria), 1) def test_build_no_limit_no_token(self):