From 9c0d604f67134f9d14633d3e58713b0d0f159fb3 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Mon, 28 Sep 2026 11:59:56 +0200 Subject: [PATCH] feat: cover the remaining api endpoints Rust (async + blocking) and Python (sync + async), with live tests in both: - functions: by_ids now calls /functions/byids instead of filtering a 10000-row listing client-side; new filter (FunctionFilter/Form) and search - timeseries: get_by_id, recommend_value_type, and listen_datapoints over the /timeseries/datapoints/listen WebSocket (token in the subprotocol) - datasets: get_by_id - resources: export_graph / import_graph, plus streaming to/from-path forms; fetch_nearest gains its first SDK-level test - tenant: features, settings_permissions, llm_settings, update_llm_settings (the SDK's first JSON PUT, via execute_put_request) execute_post_bytes_request now takes any Into so the graph import can stream from disk. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- AGENTS.md | 42 +- .../intellistream_datahub_sdk/__init__.pyi | 311 +++++++++++- .../src/datasets/async_service.rs | 16 + .../src/datasets/sync_service.rs | 11 + .../src/functions/async_service.rs | 77 +++ datahub_python_bindings/src/functions/mod.rs | 110 +++++ .../src/functions/sync_service.rs | 79 ++- datahub_python_bindings/src/lib.rs | 23 + .../src/resources/async_service.rs | 57 +++ datahub_python_bindings/src/resources/mod.rs | 78 +++ .../src/resources/sync_service.rs | 48 ++ datahub_python_bindings/src/tenant/mod.rs | 315 ++++++++++++ .../src/timeseries/async_service.rs | 57 +++ .../src/timeseries/live.rs | 277 +++++++++++ datahub_python_bindings/src/timeseries/mod.rs | 5 + .../src/timeseries/sync_service.rs | 51 ++ python_tests/test_functions.py | 31 +- python_tests/test_graph_transfer.py | 64 +++ python_tests/test_resources.py | 38 +- python_tests/test_single_reads.py | 65 +++ python_tests/test_tenant.py | 64 +++ src/blocking.rs | 57 ++- src/datasets/mod.rs | 8 + src/datasets/tests.rs | 22 + src/functions/mod.rs | 135 ++++-- src/functions/test.rs | 52 ++ src/generic.rs | 32 +- src/lib.rs | 6 + src/resources/mod.rs | 132 +++++ src/resources/tests.rs | 128 +++++ src/tenant/mod.rs | 277 +++++++++++ src/timeseries/datapoint_listen.rs | 455 ++++++++++++++++++ src/timeseries/mod.rs | 74 +++ src/timeseries/test.rs | 98 ++++ 34 files changed, 3236 insertions(+), 59 deletions(-) create mode 100644 datahub_python_bindings/src/tenant/mod.rs create mode 100644 datahub_python_bindings/src/timeseries/live.rs create mode 100644 python_tests/test_graph_transfer.py create mode 100644 python_tests/test_single_reads.py create mode 100644 python_tests/test_tenant.py create mode 100644 src/tenant/mod.rs create mode 100644 src/timeseries/datapoint_listen.rs diff --git a/AGENTS.md b/AGENTS.md index f6c68ab..f919a01 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -68,7 +68,7 @@ This crate is a thin async HTTP SDK around a DataHub-style REST API. Entry point - `time_series` (`src/timeseries/`) — `TimeSeries` + datapoint ingestion/retrieval. Neither `TimeSeries` nor `TimeSeriesUpdate` has **`securityCategories`**: it was stored, writable and returned, but nothing ever read it — no part in access control (dataset grants are Keycloak organization groups), no query filtering on it, and the backend silently dropped any id that did not already exist, so the field never round-tripped. It has been removed server-side along with its join table, and the api reads request bodies strictly, so sending it is now a 400. Files keep their own `securityCategories` (`INode` in `src/generic.rs`) — separate entity, separate question. `ListFieldU64` went with it: it was the only field of that type, so the Python wrapper class is gone too (`ListFieldStr` and `ListFieldIdCollection` remain). **`tableEngine`** went the same way, but is *not* a 400 everywhere: no read had returned it since the api marked it `@JsonIgnore` — which ClickHouse engine backs a series is an internal storage decision — and the write side is gone. `Timeseries` keeps it in a `@JsonIgnoreProperties` list so an older SDK's body is **accepted and the field ignored**; the update form (`TimeseriesFields`) has no such list, so sending it *there* is a 400. - `units` (`src/unit/`) - `events` (`src/events/`) — event CRUD, filter/search, plus the vocabulary endpoints (`list_types`/`search_types` and the same pair for sub-types, statuses and sources, over `EventDimension`). Those answer "what values does this tenant actually use" for the four categorical fields and back filter dropdowns; they read small server-side dimension tables rather than scanning events, so they are cheap but *eventually consistent* with the events. Note the route asymmetry the SDK hides: `/events/list/{plural}` but `/events/search/{singular}`. `EventUpdate` has **no `event_time` and no `external_id`**: both identify an event rather than describe it, and each was dropped from the api's update form after it had spent a while validating the field, echoing the new value back with a 200 and then failing to apply it — so sending either is now a 400 naming it. The events table is partitioned by `event_time`, so ClickHouse refuses that mutation outright. `externalId` maps to the *set* of event UUIDs behind it, events sharing one being the lifecycle of a single logical event, so a rename would take every sibling along — and when the caller identified the event by UUID the server had no old value to re-key with, leaving the "renamed" event resolvable under neither id. Re-key by creating a new event and deleting the old; record a corrected time the same way. -- `resources` (`src/resources/`) — the generic node service. Its reads span **every** node type and answer with [`Node`](#the-polymorphic-node-type) rather than one flat shape; relationship edges live in `src/relations/` (`EdgeProxy`, `RelForm`, `RelatedNode`) +- `resources` (`src/resources/`) — the generic node service. Its reads span **every** node type and answer with [`Node`](#the-polymorphic-node-type) rather than one flat shape; relationship edges live in `src/relations/` (`EdgeProxy`, `RelForm`, `RelatedNode`). `export_graph`/`import_graph` (and the `_to_path`/`_from_path` streaming pair) move a whole connected component as a gzip file keyed by external id; import skips what exists, so re-importing into the source tenant is a no-op. The export walks the graph projection, which lags the write — a component exported straight after creating it comes back empty, so wait on `fetch_related` first. - `edges` (`src/relations/service.rs`) — the `/edges` endpoints: `get`/`by_ids`/`create`/`delete` plus the relationship-type catalogue (`types`/`create_types`). Edges normally come into being through `resources.create(nodes, relations)`; this service is for linking resources that already exist and for reading or deleting an edge on its own. `get` answers an unknown id with 404 and a `problem+json` body; `by_ids`, like every batch lookup, answers 200 with the found subset and silently omits what is missing. (`get` used to be 200-and-nothing despite documenting a 404 — api #275 made single-resource by-id GETs consistently 404 and deliberately left batch lookups alone.) Two further behaviours are worth knowing, and are documented at each call site: - `create_types` answers a duplicate name with **409** `duplicate`, naming `name` in `fields`. It used to fail silently — the unique-hash collision surfaced at commit, after the handler returned, so the caller got a 200 with an empty *body*. `test_duplicate_relationship_type_conflicts` encoded the intended 409 while that was true and is a regression guard now. Still true, and still worth knowing: the service has **no find-or-create**, and a batch is one transaction, so a single duplicate rolls back the valid new types alongside it. @@ -87,14 +87,20 @@ This crate is a thin async HTTP SDK around a DataHub-style REST API. Entry point - `datasets` (`src/datasets/`) - `files` (`src/files/`) — raw-`PUT` upload via `execute_file_upload_request` (content is the body, metadata rides in `X-Datahub-*` headers), plus directory listing, get/search, `FileUpdate` (rename/move/re-dataset), trash + restore, delete, and download (`download` in memory, `download_to_path` streamed) - `subscriptions` (`src/subscriptions/`) — subscription CRUD, plus `listen.rs`: WebSocket listening against the api's subscription-listen endpoint (`tokio-tungstenite`). Reads follow the same split as every other collection: `list(limit)` over `GET /subscriptions?limit=` and `filter(form)` over `POST /subscriptions/filter`. Both are recent — `POST /subscriptions/list` was subscriptions-only (a `limit` defaulting to 100 where the api defaulted to 1000, an unvalidated sort property, and no cursor, so a tenant past one page could not reach the rest) and the api removed it rather than aliasing it, so a client that has not moved gets a 404. `SubscriptionFilterForm` has the family shape (`filter`, `limit: Option`, flattened `PageRequest`), and `SubscriptionFilter` carries everything the api's does: `id`, `externalId`, `name`, `timeseries`, `createdTime`, `lastUpdatedTime`. It is deliberately not a `NodeFilter` — a subscription has no `source`, `labels` or `metadata` — and note the criteria are `createdTime`/`lastUpdatedTime` while the entity spells its timestamps `dateCreated`/`lastUpdated`; that asymmetry is the api's. In Python, `filter()` takes the criteria as keywords or `filter=` and returns a `Page`, like `datasets.filter`. -- `functions` (`src/functions/`) — `create`, `list`, `get_by_id`, `update`, `delete`, plus a - client-side `by_ids`/`by_external_id`. The api **does** serve `/functions/byids`, `/filter` and - `/search` (platform #131) — **the SDK has not wired them yet**, so `by_ids` still lists and - filters locally, asking for the largest page the api allows; a tenant past 10000 functions - silently misses its oldest. Wiring the three is the fix, not a bigger page. +- `functions` (`src/functions/`) — `create`, `list`, `get_by_id`, `by_ids`, `filter`, `search`, + `update`, `delete`, plus `by_external_id` over `by_ids`. `filter` takes `FunctionFilterForm` + around `FunctionFilter` — the shared `NodeFilter` plus `dataSetId`, which is all the api's + `FunctionFilter` declares. `by_ids` used to list and filter client-side, so a tenant past 10000 + functions silently lost its oldest; it calls `/functions/byids` now. `Function::related_resources` is **always empty**: `FunctionTransformer` never joins the edges in, and unlike `/resources/create` the create echo is no exception, because it re-reads the rows through that same transformer. `name` is `Option` here but non-null on the api. +- `tenant` (`src/tenant/`) — `features`, `settings_permissions`, `llm_settings` and + `update_llm_settings`. Every answer is a bare object, not an `items` envelope, so each type has + its own `DataWrapperDeserialization`. `update_llm_settings` is a `PUT` that **replaces** the + object — a field left `None` is cleared — except `api_key`, where `None` keeps the stored + credential; `TenantLlmSettingsForm::from(&stored)` builds the form that saves settings unchanged. + This is the SDK's only JSON `PUT`, via `execute_put_request`. - `labels` (`src/labels/`) — label CRUD (`list`/`get`/`create`/`update`/`delete`). Note the entity type is `labels::Label`, deliberately *not* re-exported at the crate root because `resources::*` already brings a different graph-DTO `Label` there. ### The polymorphic node type (`src/nodes.rs`) @@ -169,6 +175,25 @@ Synchronous mirror of the async API behind the `blocking` cargo feature — the When a datapoint/event send can't get through, ingestion spools to a segmented, zstd-compressed NDJSON log on disk and flushes automatically on a later ingest call. Invariants to preserve: memory use is bounded by a single segment (plain append-only active segment, zstd-sealed at ~50 MiB rollover via temp file + atomic rename, drained oldest-first one segment at a time); bounded by time retention (whole segments past the window dropped, expired records skipped on read) and a size cap (oldest segment deleted); a torn trailing line from an unclean shutdown is skipped on read. Each on-disk line is `\t`; the spool is content-agnostic. +### Live datapoint tail (`src/timeseries/datapoint_listen.rs`) + +`TimeSeriesService::listen_datapoints` opens `ws(s):///timeseries/datapoints/listen`, a +non-durable tail from *latest* with no subscription entity behind it and nothing to ack — the +at-least-once path is still `SubscriptionsService::listen`. Two things differ from that listener: + +- **The token rides in `Sec-WebSocket-Protocol`**, as `datahub.bearer.,datahub.v1`, not in an + `Authorization` header (the api reads the subprotocol so browsers can use the same socket). **No + space after the comma**: tungstenite splits the offer on `,` without trimming and would reject + the echoed `datahub.v1`. `the_token_rides_in_the_subprotocol_and_points_arrive` fails if the + space comes back. The api switched to this on 2026-09-25 (platform `0a9161be`); an api started + before that answers "Server sent no subprotocol", which means the api is stale, not the SDK. +- **The api authenticates after the 101**, so a bad token or missing role arrives as a 1008 + policy close. `next` returns that as `ListenError::Handshake` instead of reconnecting into the + same refusal. + +`recommend_value_type` sits beside it: a bare `ValueTypeRecommendation`, and an unknown unit is +the generic default with `recognized: false`, not an error. + ### Binary datapoint ingest (`src/timeseries/binary.rs`) `TimeSeriesService::insert_datapoints_binary` is the second ingest path, to `POST /timeseries/data/binary`: the same `DatapointsCollection` input, resolved through `/timeseries/byids` (cached per service instance; needs read access on the dataset), checked against each series' value type locally, cut into Arrow IPC frames at the contract's caps, zstd-compressed per frame (mandatory: level 1, 3 or 9, default 9) and posted through `execute_post_bytes_request`. Every refusal on this path is a problem document typed for its status — `invalid-frame` (400), `unknown-timeseries` (404), `request-too-large` (413), `unsupported-media-type` (415), `value-type-mismatch` and `external-id-mismatch` (422), `too-many-in-flight` (429) — with the kebab-case sub-case in a `reason` extension beside it. It was one type, `datapoint-block-rejected`, answering with all six statuses until the api split it (platform #120); the SDK matches the new slugs only, and nothing here recognises the old one — the binary path has never been in a release, so no caller can be on the other side of that change. `is_stale_series_rejection` reads the slug to decide the one re-resolve-and-retry, and reads it from the problem rather than the body text because the SDK's own pre-flight 404 names the series it could not resolve — an external id spelling `unknown-timeseries` used to trigger the retry. `FrameWriter` builds one frame and is public; `cut_into_writers` and `pack_requests` hold the caps. The byte layout is the platform's `binary_datapoints_format.md`. The arrow-rs crates (`arrow-array`, `arrow-schema`, `arrow-ipc`) exist for this path and are the seed of the Arrow read path. Not in the Python bindings yet. The ignored `timeseries::tests::test_datapoints_binary` is the live twin of `test_datapoints` and needs a backend that serves the endpoint; the writer's own tests in `binary.rs` run offline. @@ -369,9 +394,8 @@ collection that had them. and `python_tests/test_plain_listings.py`. - **`resources.list` is typed like every other `/resources` read** — it spans all six node types and answers each row as its own [`Node`](#the-polymorphic-node-type) variant. -- **The listing is not a way to fetch everything.** `FunctionsService::by_ids` does not yet call - the api's `/functions/byids` and filters a listing client-side; that listing used to be uncapped, - so it now asks for 10000 and a tenant past that silently misses its oldest functions. +- **The listing is not a way to fetch everything.** A cap of 10000 truncates; narrow with + `filter` or look entities up with `by_ids`. Names that are gone rather than aliased, the way the filter refactors handled theirs: `POST /datasets/list` (a stale caller gets **405** — `GET /datasets/{id}` matches the path and diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index c88c205..101719d 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -190,6 +190,8 @@ class DataHubClient: def labels(self) -> LabelsServiceSync: ... @property def edges(self) -> EdgesServiceSync: ... + @property + def tenant(self) -> TenantServiceSync: ... class AsyncDataHubClient: @@ -243,6 +245,8 @@ class AsyncDataHubClient: def edges(self) -> EdgesServiceAsync: ... @property def datasets(self) -> DatasetsServiceAsync: ... + @property + def tenant(self) -> TenantServiceAsync: ... # ====================== Identifiers & search ====================== @@ -611,8 +615,76 @@ class RetrieveFilter: def cursor(self) -> str | None: ... +class LiveDatapoint: + """One point from ``timeseries.listen_datapoints()``. ``value`` is a string for every + value type.""" + @property + def external_id(self) -> str: ... + @property + def value_type(self) -> str | None: ... + @property + def timestamp(self) -> str: ... + @property + def value(self) -> str: ... + + +class ValueTypeRecommendation: + @property + def unit_external_id(self) -> str: ... + @property + def recommended_value_type(self) -> str: + """One of ``BIGINT``, ``FLOAT``, ``FLOAT32``, ``NUMERIC``, ``DECIMAL32``, ``TEXT``, + ``MIXED``.""" + @property + def reason(self) -> str: ... + @property + def recognized(self) -> bool: + """``False`` when the unit matched nothing specific and the generic default came back.""" + + +class DatapointListener: + """A live tail of datapoints. ``for point in listener:`` blocks until the next one. + + Nothing is durable: points written while the connection is down are not replayed. A dropped + connection is re-established transparently; a refusal (bad token, missing role, connection + limit) raises instead.""" + def __iter__(self) -> DatapointListener: ... + def __next__(self) -> LiveDatapoint: ... + def next_datapoint(self) -> LiveDatapoint | None: ... + def subscribe(self, external_ids: list[str]) -> None: + """Add timeseries. Ids you cannot read are dropped silently, not reported.""" + def unsubscribe(self, external_ids: list[str]) -> None: ... + def set_timeseries(self, external_ids: list[str]) -> None: + """Replace the whole live set.""" + def close(self) -> None: ... + def __enter__(self) -> DatapointListener: ... + def __exit__(self, exc_type: Any, exc_value: Any, traceback: Any) -> None: ... + + +class DatapointListenerAsync: + def __aiter__(self) -> DatapointListenerAsync: ... + async def __anext__(self) -> LiveDatapoint: ... + async def next_datapoint(self) -> LiveDatapoint | None: ... + async def subscribe(self, external_ids: list[str]) -> None: ... + async def unsubscribe(self, external_ids: list[str]) -> None: ... + async def set_timeseries(self, external_ids: list[str]) -> None: ... + async def close(self) -> None: ... + async def __aenter__(self) -> DatapointListenerAsync: ... + async def __aexit__(self, exc_type: Any, exc_value: Any, traceback: Any) -> None: ... + + class TimeSeriesServiceSync: def list(self, limit: int | None = None) -> list[TimeSeries]: ... + def get_by_id(self, id: int) -> TimeSeries | None: + """One series by numeric id; raises on 404. A 404 does not tell you the id is free — a + series you may not read is reported as missing rather than forbidden.""" + def recommend_value_type(self, unit_external_id: str) -> ValueTypeRecommendation: + """The value type that compresses best for a unit while representing it faithfully. + Advice, not a constraint. An unknown unit is not an error: it answers the generic default + with ``recognized=False``.""" + def listen_datapoints(self, external_ids: list[str] = ...) -> DatapointListener: + """Open a live tail of the datapoints written to these timeseries, starting at the latest + point. For at-least-once delivery use ``subscriptions.listen`` instead.""" def create(self, input: list[TimeSeries]) -> list[TimeSeries]: ... def by_ids(self, input: list[Identifiable]) -> list[TimeSeries]: ... def delete(self, input: list[Identifiable]) -> None: ... @@ -684,6 +756,9 @@ class TimeSeriesServiceSync: class TimeSeriesServiceAsync: async def list(self, limit: int | None = None) -> list[TimeSeries]: ... + async def get_by_id(self, id: int) -> TimeSeries | None: ... + async def recommend_value_type(self, unit_external_id: str) -> ValueTypeRecommendation: ... + async def listen_datapoints(self, external_ids: list[str] = ...) -> DatapointListenerAsync: ... async def create(self, input: list[TimeSeries]) -> list[TimeSeries]: ... async def by_ids(self, input: list[Identifiable]) -> list[TimeSeries]: ... async def delete(self, input: list[Identifiable]) -> None: ... @@ -1152,6 +1227,8 @@ class DatasetUpdate: # `update(...)`: there is no write_protected/deactivated — both were removed server-side as inert. class DatasetsServiceSync: def list(self, limit: int | None = None) -> list[Dataset]: ... + def get_by_id(self, id: int) -> Dataset | None: + """One data set by numeric id; raises on 404.""" def create(self, input: list[Dataset]) -> list[Dataset]: ... def by_ids(self, input: list[Identifiable]) -> list[Dataset]: ... def delete(self, input: list[Identifiable]) -> None: ... @@ -1196,6 +1273,7 @@ class DatasetsServiceSync: class DatasetsServiceAsync: async def list(self, limit: int | None = None) -> list[Dataset]: ... + async def get_by_id(self, id: int) -> Dataset | None: ... async def create(self, input: list[Dataset]) -> list[Dataset]: ... async def by_ids(self, input: list[Identifiable]) -> list[Dataset]: ... async def delete(self, input: list[Identifiable]) -> None: ... @@ -1656,8 +1734,49 @@ class ResourceFilter: ) -> None: ... +class GraphImportResult: + """What ``resources.import_graph()`` did.""" + @property + def nodes_created(self) -> int: ... + @property + def relations_created(self) -> int: ... + @property + def nodes_skipped_existing(self) -> int: + """Skipped because a node with the same external id already exists.""" + @property + def nodes_skipped_timeseries(self) -> list[str]: + """Timeseries in the file that do not exist here. They cannot be created through the + resource api — create them through ``timeseries`` first, then import again.""" + @property + def relations_skipped(self) -> int: ... + @property + def data_set_references_dropped(self) -> int: ... + @property + def segments(self) -> int: + """Transactions committed; each segment of 50,000 objects is atomic on its own.""" + @property + def warnings(self) -> list[PolicyWarning]: ... + + class ResourcesServiceSync: def list(self, limit: int | None = None) -> list[Node]: ... + def export_graph(self, id: int) -> bytes: + """The whole connected graph component around one resource, as a gzip-compressed file. + It names everything by external id, so it imports into another tenant with + ``import_graph``. Over 2,000,000 nodes or relationships raises a 400 rather than exporting + part of it. Use ``export_graph_to_path`` for anything large.""" + def export_graph_to_path(self, id: int, destination: str) -> int: + """``export_graph``, streamed to ``destination``. Returns the number of bytes written.""" + def import_graph(self, file: bytes) -> GraphImportResult: + """Recreate the resources and relationships of an ``export_graph`` file. + + What already exists is skipped — nodes by external id, relationships by (from, to, type) — + so re-importing into the source tenant is a no-op, and after a failure the same file can + simply be sent again. Over 512 MB, or 2,000,000 nodes or relationships, raises a 413. + The graph projection lags writes, so export a component only once it is readable through + ``fetch_related``.""" + def import_graph_from_path(self, source: str) -> GraphImportResult: + """``import_graph``, streaming the file from disk.""" def create( self, nodes: list[Node], relations: list[RelForm] | None = None ) -> GraphResult: ... @@ -1713,6 +1832,10 @@ class ResourcesServiceSync: class ResourcesServiceAsync: async def list(self, limit: int | None = None) -> list[Node]: ... + async def export_graph(self, id: int) -> bytes: ... + async def export_graph_to_path(self, id: int, destination: str) -> int: ... + async def import_graph(self, file: bytes) -> GraphImportResult: ... + async def import_graph_from_path(self, source: str) -> GraphImportResult: ... async def create( self, nodes: list[Node], relations: list[RelForm] | None = None ) -> GraphResult: ... @@ -2304,6 +2427,24 @@ class Function: FunctionIdentifiable = Union[Function, IdCollection, int, str] +class FunctionFilter: + """Criteria for ``functions.filter()`` and the ``filter`` of ``functions.search()``: the + criteria every node type shares, plus ``data_set_id``. ``data_set_id=None`` places no + restriction; ``[]`` matches nothing.""" + def __init__( + self, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + ) -> None: ... + + class FunctionsServiceSync: def create(self, input: list[Function]) -> list[Function]: ... def list(self, limit: int | None = None) -> list[Function]: ... @@ -2311,10 +2452,38 @@ class FunctionsServiceSync: """One function by numeric id; raises on 404. A 404 does not tell you the id is free — a function you may not read is reported as - missing rather than forbidden. Prefer this to ``by_ids`` when you have the id: functions - have no ``/byids`` endpoint, so ``by_ids`` pages the listing and filters client-side. + missing rather than forbidden. """ - def by_ids(self, input: list[FunctionIdentifiable]) -> list[Function]: ... + def by_ids(self, input: list[FunctionIdentifiable]) -> list[Function]: + """What does not exist, or you may not read, is left out rather than raising.""" + def filter( + self, + *, + filter: FunctionFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Newest created first unless ``sort_by`` says otherwise.""" + def search( + self, + query: str, + filter: FunctionFilter | None = None, + limit: int | None = None, + ) -> list[Function]: + """Free-text search, best match first. ``query`` is 3–140 characters; ``limit`` defaults + to 100 and caps at 1000.""" def by_external_id(self, external_id: str) -> Function: ... def update(self, input: list[ResourceUpdate]) -> GraphResult: """Update functions in place; ``geolocation`` is ignored, being asset-only. @@ -2329,6 +2498,30 @@ class FunctionsServiceAsync: async def list(self, limit: int | None = None) -> list[Function]: ... async def get_by_id(self, id: int) -> Function | None: ... async def by_ids(self, input: list[FunctionIdentifiable]) -> list[Function]: ... + async def filter( + self, + *, + filter: FunctionFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: ... + async def search( + self, + query: str, + filter: FunctionFilter | None = None, + limit: int | None = None, + ) -> list[Function]: ... async def by_external_id(self, external_id: str) -> Function: ... async def update(self, input: list[ResourceUpdate]) -> GraphResult: ... async def delete(self, input: list[FunctionIdentifiable]) -> None: ... @@ -2504,3 +2697,115 @@ class EdgesServiceAsync: async def delete(self, input: list[EdgeIdentifiable]) -> None: ... async def types(self) -> list[RelationshipType]: ... async def create_types(self, input: list[RelTypeForm]) -> list[RelationshipType]: ... + + +# ====================== Policies ====================== + +class PolicyWarning: + """A naming-policy violation that was allowed through and recorded for review.""" + @property + def index(self) -> int: ... + @property + def external_id(self) -> str: ... + @property + def policy(self) -> str | None: ... + @property + def message(self) -> str | None: ... + @property + def suggestion(self) -> str | None: ... + + +# ====================== Tenant ====================== + +class TenantFeatures: + """Which optional features are enabled for your tenant. A disabled feature's endpoints may + still exist and answer 404 or 403.""" + @property + def files(self) -> bool: ... + @property + def policy(self) -> bool: ... + @property + def streaming(self) -> bool: ... + @property + def chat(self) -> bool: ... + + +class SettingsPermission: + @property + def read(self) -> bool: ... + @property + def write(self) -> bool: ... + + +class TenantLlmSettings: + """The model your organization's assistant runs on. The API key is never returned.""" + @property + def provider(self) -> str | None: + """``"anthropic"`` or ``"openai-compatible"``.""" + @property + def model(self) -> str | None: ... + @property + def base_url(self) -> str | None: ... + @property + def reasoning_effort(self) -> str | None: ... + @property + def effort(self) -> str | None: + """One of ``low``, ``medium``, ``high``, ``xhigh``, ``max``.""" + @property + def turn_timeout(self) -> str | None: ... + @property + def max_output_tokens(self) -> int | None: ... + @property + def max_iterations(self) -> int | None: ... + @property + def instructions(self) -> str | None: ... + @property + def api_key_set(self) -> bool: + """Whether a credential is stored.""" + @property + def configured(self) -> bool: + """Whether this amounts to a model that can actually be called.""" + + +class TenantServiceSync: + def features(self) -> TenantFeatures: ... + def settings_permissions(self) -> dict[str, SettingsPermission]: + """What you may read and write, per settings scope (``"llm"``, …). For gating a UI; the + settings calls enforce the same grants.""" + def llm_settings(self) -> TenantLlmSettings: + """Needs the ``llm`` read grant (403 without).""" + def update_llm_settings( + self, + provider: str | None = None, + model: str | None = None, + api_key: str | None = None, + base_url: str | None = None, + reasoning_effort: str | None = None, + effort: str | None = None, + turn_timeout: str | None = None, + max_output_tokens: int | None = None, + max_iterations: int | None = None, + instructions: str | None = None, + ) -> TenantLlmSettings: + """**Replace** the model configuration: an argument left out is cleared. The exception is + ``api_key``, where ``None`` or empty keeps the stored credential. Needs the ``llm`` write + grant.""" + + +class TenantServiceAsync: + async def features(self) -> TenantFeatures: ... + async def settings_permissions(self) -> dict[str, SettingsPermission]: ... + async def llm_settings(self) -> TenantLlmSettings: ... + async def update_llm_settings( + self, + provider: str | None = None, + model: str | None = None, + api_key: str | None = None, + base_url: str | None = None, + reasoning_effort: str | None = None, + effort: str | None = None, + turn_timeout: str | None = None, + max_output_tokens: int | None = None, + max_iterations: int | None = None, + instructions: str | None = None, + ) -> TenantLlmSettings: ... diff --git a/datahub_python_bindings/src/datasets/async_service.rs b/datahub_python_bindings/src/datasets/async_service.rs index 733b2e4..acbb3f0 100644 --- a/datahub_python_bindings/src/datasets/async_service.rs +++ b/datahub_python_bindings/src/datasets/async_service.rs @@ -36,6 +36,22 @@ impl PyDatasetsServiceAsync { }) } + /// One data set by numeric id. Raises on 404. + fn get_by_id<'p>(&self, py: Python<'p>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + let result = service + .datasets + .get_by_id(id) + .await + .map_err(|e| crate::datahub_err(e))?; + Ok(result + .get_items() + .first() + .map(|d| PyDataset::with_client(d.clone(), service.clone()))) + }) + } + fn create<'p>(&self, py: Python<'p>, input: Vec) -> PyResult> { let datasets: Vec = input.iter().cloned().map(Dataset::from).collect(); let service = self.api_service.clone(); diff --git a/datahub_python_bindings/src/datasets/sync_service.rs b/datahub_python_bindings/src/datasets/sync_service.rs index ac65c2d..a7f7ccb 100644 --- a/datahub_python_bindings/src/datasets/sync_service.rs +++ b/datahub_python_bindings/src/datasets/sync_service.rs @@ -89,6 +89,17 @@ impl PyDatasetsServiceSync { }) } + /// One data set by numeric id. Raises on 404. + fn get_by_id(&self, py: Python<'_>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + let result = py.detach(|| self.runtime.block_on(service.datasets.get_by_id(id))); + let result = result.map_err(crate::datahub_err)?; + Ok(result + .get_items() + .first() + .map(|d| PyDataset::with_client(d.clone(), service.clone()))) + } + /// Datasets matching every criterion on the filter, newest first. #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, labels=None, metadata=None, created_time=None, last_updated_time=None, diff --git a/datahub_python_bindings/src/functions/async_service.rs b/datahub_python_bindings/src/functions/async_service.rs index f24a2c4..b1982a3 100644 --- a/datahub_python_bindings/src/functions/async_service.rs +++ b/datahub_python_bindings/src/functions/async_service.rs @@ -56,6 +56,83 @@ impl PyFunctionsServiceAsync { }) } + /// Functions matching every criterion, newest created first unless `sort_by` says otherwise. + /// Returns a `Page`; send its `next_cursor` back as `cursor`, with the same sort, for the next. + /// + /// Pass either `filter=` or the individual keywords, not both. + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + data_set_id=None, limit=None, sort_by=None, sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] + fn filter<'py>( + &self, + py: Python<'py>, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult> { + let form = crate::functions::build_function_filter_form( + filter, id, external_id, name, source, labels, metadata, created_time, + last_updated_time, data_set_id, limit, sort_by, sort_order, cursor, + )?; + let service = self.api_service.clone(); + future_into_py(py, async move { + let result = service + .functions + .filter(&form) + .await + .map_err(|e| crate::datahub_err(e))?; + let next_cursor = result.next_cursor().map(str::to_string); + let items: Vec = result + .get_items() + .iter() + .cloned() + .map(|f| PyFunction::with_client(f, service.clone())) + .collect(); + Python::attach(|py| crate::PyPage::new(py, items, next_cursor)) + }) + } + + /// Free-text search over functions, best match first. The phrase selects and `filter` only + /// removes. `query` is required at 3–140 characters; `limit` defaults to 100 and caps at 1000. + #[pyo3(signature = (query, filter = None, limit = None))] + fn search<'py>( + &self, + py: Python<'py>, + query: String, + filter: Option, + limit: Option, + ) -> PyResult> { + let form = crate::search_form(query, filter.map(|f| f.inner), limit); + let service = self.api_service.clone(); + future_into_py(py, async move { + let result = service + .functions + .search(&form) + .await + .map_err(|e| crate::datahub_err(e))?; + Ok(result + .get_items() + .iter() + .cloned() + .map(|f| PyFunction::with_client(f, service.clone())) + .collect::>()) + }) + } + + /// Look up functions by id or external id. What does not exist, or you may not read, is + /// left out rather than raising. fn by_ids<'py>( &self, py: Python<'py>, diff --git a/datahub_python_bindings/src/functions/mod.rs b/datahub_python_bindings/src/functions/mod.rs index 3b11039..468e11e 100644 --- a/datahub_python_bindings/src/functions/mod.rs +++ b/datahub_python_bindings/src/functions/mod.rs @@ -283,6 +283,7 @@ pub(crate) fn json_to_py<'py>( pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; + m.add_class::()?; m.add_class::()?; m.add_class::()?; Ok(()) @@ -402,3 +403,112 @@ impl PyFunction { filter } } + +/// Criteria for `functions.filter()` and the `filter` of `functions.search()`: the criteria every +/// node type shares, plus `data_set_id`. Same rules as `ResourceFilter` — patterns, AND across +/// fields, and `data_set_id=None` (no restriction) differing from `[]` (matches nothing). +#[pyclass(module = "intellistream_datahub_sdk", name = "FunctionFilter", from_py_object)] +#[derive(Clone)] +pub struct PyFunctionFilter { + pub inner: intellistream_datahub_sdk::functions::FunctionFilter, +} + +#[pymethods] +impl PyFunctionFilter { + #[new] + #[pyo3(signature = (id=None, external_id=None, name=None, source=None, labels=None, + metadata=None, created_time=None, last_updated_time=None, + data_set_id=None))] + #[allow(clippy::too_many_arguments)] + fn new( + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + ) -> Self { + Self { + inner: build_function_filter( + id, + external_id, + name, + source, + labels, + metadata, + created_time, + last_updated_time, + data_set_id, + ), + } + } +} + +#[allow(clippy::too_many_arguments)] +pub(crate) fn build_function_filter( + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, +) -> intellistream_datahub_sdk::functions::FunctionFilter { + intellistream_datahub_sdk::functions::FunctionFilter { + node: intellistream_datahub_sdk::filters::NodeFilter { + id, + external_id: crate::opt_patterns(external_id), + name: crate::opt_patterns(name), + source: crate::opt_patterns(source), + labels: crate::opt_patterns(labels), + metadata, + created_time: created_time.map(Into::into), + last_updated_time: last_updated_time.map(Into::into), + }, + data_set_id: crate::opt_data_set_refs(data_set_id), + } +} + +/// Shared by the sync and async `filter` bindings. +#[allow(clippy::too_many_arguments)] +pub(crate) fn build_function_filter_form( + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, +) -> PyResult { + let any_keyword = id.is_some() + || external_id.is_some() + || name.is_some() + || source.is_some() + || labels.is_some() + || metadata.is_some() + || created_time.is_some() + || last_updated_time.is_some() + || data_set_id.is_some(); + let from_keywords = build_function_filter( + id, external_id, name, source, labels, metadata, created_time, last_updated_time, + data_set_id, + ); + let filter = crate::resolve_filter(filter.map(|f| f.inner), from_keywords, any_keyword)?; + let mut form = intellistream_datahub_sdk::functions::FunctionFilterForm::new(filter); + if let Some(limit) = limit { + form = form.with_limit(limit); + } + Ok(form.with_paging(crate::build_page_request(sort_by, sort_order, cursor))) +} diff --git a/datahub_python_bindings/src/functions/sync_service.rs b/datahub_python_bindings/src/functions/sync_service.rs index 966c3c8..583e29f 100644 --- a/datahub_python_bindings/src/functions/sync_service.rs +++ b/datahub_python_bindings/src/functions/sync_service.rs @@ -50,6 +50,82 @@ impl PyFunctionsServiceSync { }) } + /// Functions matching every criterion, newest created first unless `sort_by` says otherwise. + /// Returns a `Page`; send its `next_cursor` back as `cursor`, with the same sort, for the next. + /// + /// Pass either `filter=` or the individual keywords, not both. + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + data_set_id=None, limit=None, sort_by=None, sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] + fn filter( + &self, + py: Python<'_>, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult { + let form = crate::functions::build_function_filter_form( + filter, id, external_id, name, source, labels, metadata, created_time, + last_updated_time, data_set_id, limit, sort_by, sort_order, cursor, + )?; + let service = self.api_service.clone(); + let (items, next_cursor) = py.detach(|| { + let result = self + .runtime + .block_on(service.functions.filter(&form)) + .map_err(|e| crate::datahub_err(e))?; + let next_cursor = result.next_cursor().map(str::to_string); + let items: Vec = result + .get_items() + .iter() + .cloned() + .map(|f| PyFunction::with_client(f, service.clone())) + .collect(); + Ok::<_, pyo3::PyErr>((items, next_cursor)) + })?; + crate::PyPage::new(py, items, next_cursor) + } + + /// Free-text search over functions, best match first. The phrase selects and `filter` only + /// removes. `query` is required at 3–140 characters; `limit` defaults to 100 and caps at 1000. + #[pyo3(signature = (query, filter = None, limit = None))] + fn search( + &self, + py: Python<'_>, + query: String, + filter: Option, + limit: Option, + ) -> PyResult> { + let form = crate::search_form(query, filter.map(|f| f.inner), limit); + let service = self.api_service.clone(); + py.detach(|| { + let result = self + .runtime + .block_on(service.functions.search(&form)) + .map_err(|e| crate::datahub_err(e))?; + Ok(result + .get_items() + .iter() + .cloned() + .map(|f| PyFunction::with_client(f, service.clone())) + .collect()) + }) + } + + /// Look up functions by id or external id. What does not exist, or you may not read, is + /// left out rather than raising. fn by_ids( &self, py: Python<'_>, @@ -87,8 +163,7 @@ impl PyFunctionsServiceSync { /// One function by numeric id. /// /// Raises on 404 — and a 404 does not tell you the id is free: a function you may not read is - /// reported as missing rather than forbidden. Prefer this to `by_ids` when you have the id: - /// `by_ids` has no endpoint behind it and pages the whole listing to filter client-side. + /// reported as missing rather than forbidden. fn get_by_id(&self, py: Python<'_>, id: u64) -> PyResult> { let service = self.api_service.clone(); let result = py.detach(|| self.runtime.block_on(service.functions.get_by_id(id))); diff --git a/datahub_python_bindings/src/lib.rs b/datahub_python_bindings/src/lib.rs index a9a81ba..9e32c26 100644 --- a/datahub_python_bindings/src/lib.rs +++ b/datahub_python_bindings/src/lib.rs @@ -11,6 +11,7 @@ mod subscriptions; pub mod timeseries; pub mod units; mod functions; +mod tenant; use crate::datasets::PyDataset; use crate::datasets::async_service::PyDatasetsServiceAsync; @@ -406,6 +407,16 @@ impl PySyncClient { } } + + + #[getter] + fn tenant(&self) -> crate::tenant::PyTenantServiceSync { + crate::tenant::PyTenantServiceSync { + api_service: self.inner.clone(), + runtime: self.runtime.clone(), + } + } + } #[pyclass(module = "intellistream_datahub_sdk", name = "AsyncDataHubClient")] @@ -563,6 +574,15 @@ impl PyAsyncClient { } } + + + #[getter] + fn tenant(&self) -> crate::tenant::PyTenantServiceAsync { + crate::tenant::PyTenantServiceAsync { + api_service: self.inner.clone(), + } + } + #[getter] fn datasets(&self) -> PyDatasetsServiceAsync { PyDatasetsServiceAsync { @@ -1291,5 +1311,8 @@ fn _core(m: &Bound<'_, PyModule>) -> PyResult<()> { assets::register(m)?; relations::register(m)?; nodes::register(m)?; + tenant::register(m)?; + m.add_class::()?; + m.add_class::()?; Ok(()) } diff --git a/datahub_python_bindings/src/resources/async_service.rs b/datahub_python_bindings/src/resources/async_service.rs index 93d2dc1..04b5fbb 100644 --- a/datahub_python_bindings/src/resources/async_service.rs +++ b/datahub_python_bindings/src/resources/async_service.rs @@ -235,6 +235,63 @@ impl PyResourcesServiceAsync { }) } + /// The whole connected graph component around one resource, as a gzip-compressed file + /// (`bytes`). It names everything by external id, so it imports into another tenant with + /// `import_graph`. Use `export_graph_to_path` for anything large. + fn export_graph<'py>(&self, py: Python<'py>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + let file = service + .resources + .export_graph(id) + .await + .map_err(|e| crate::datahub_err(e))?; + Ok(Python::attach(|py| pyo3::types::PyBytes::new(py, &file).unbind())) + }) + } + + /// `export_graph`, streamed to `destination` without holding it in memory. Returns the number + /// of bytes written. + fn export_graph_to_path<'py>(&self, py: Python<'py>, id: u64, destination: String) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .resources + .export_graph_to_path(id, destination) + .await + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// Recreate the resources and relationships of an `export_graph` file. What already exists is + /// skipped, so re-importing into the source tenant is a no-op and a failed import can simply + /// be sent again. Timeseries are not created; missing ones are listed on the result. + fn import_graph<'py>(&self, py: Python<'py>, file: &[u8]) -> PyResult> { + let service = self.api_service.clone(); + let file = file.to_vec(); + future_into_py(py, async move { + service + .resources + .import_graph(file) + .await + .map(crate::resources::PyGraphImportResult::from) + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// `import_graph`, streaming the file from disk. + fn import_graph_from_path<'py>(&self, py: Python<'py>, source: String) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .resources + .import_graph_from_path(source) + .await + .map(crate::resources::PyGraphImportResult::from) + .map_err(|e| crate::datahub_err(e)) + }) + } + /// `POST /resources/fetch-nearest` — the closest `limit` nodes carrying one of `end_labels`, /// plus the sub-graph connecting them back to the start. Starts from a numeric `id` only. #[pyo3(signature = (id, end_labels=None, limit=None, relationship_types=None, excluded_labels=None))] diff --git a/datahub_python_bindings/src/resources/mod.rs b/datahub_python_bindings/src/resources/mod.rs index 33cc3ff..30b6b60 100644 --- a/datahub_python_bindings/src/resources/mod.rs +++ b/datahub_python_bindings/src/resources/mod.rs @@ -604,3 +604,81 @@ impl PyResourceFilter { } } } + +/// Answer of `resources.import_graph()`. +#[pyclass(module = "intellistream_datahub_sdk", name = "GraphImportResult", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyGraphImportResult { + pub nodes_created: u64, + pub relations_created: u64, + /// Skipped because a node with the same external id already exists. + pub nodes_skipped_existing: u64, + /// Timeseries in the file that do not exist here; create them through the timeseries api. + pub nodes_skipped_timeseries: Vec, + pub relations_skipped: u64, + pub data_set_references_dropped: u64, + /// Transactions committed; each segment of 50,000 objects is atomic on its own. + pub segments: u64, + /// Naming-policy violations allowed through and recorded for review. + pub warnings: Vec, +} + +impl From for PyGraphImportResult { + fn from(r: intellistream_datahub_sdk::resources::GraphImportResult) -> Self { + Self { + nodes_created: r.nodes_created, + relations_created: r.relations_created, + nodes_skipped_existing: r.nodes_skipped_existing, + nodes_skipped_timeseries: r.nodes_skipped_timeseries, + relations_skipped: r.relations_skipped, + data_set_references_dropped: r.data_set_references_dropped, + segments: r.segments, + warnings: r.warnings.into_iter().map(Into::into).collect(), + } + } +} + +#[pymethods] +impl PyGraphImportResult { + fn __repr__(&self) -> String { + format!( + "GraphImportResult(nodes_created={}, relations_created={}, nodes_skipped_existing={}, relations_skipped={}, segments={})", + self.nodes_created, + self.relations_created, + self.nodes_skipped_existing, + self.relations_skipped, + self.segments + ) + } +} + +/// A naming-policy violation that was allowed through and recorded for review. +#[pyclass(module = "intellistream_datahub_sdk", name = "PolicyWarning", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyPolicyWarning { + /// Position of the offending item in the submitted batch. + pub index: u32, + pub external_id: String, + pub policy: Option, + pub message: Option, + pub suggestion: Option, +} + +impl From for PyPolicyWarning { + fn from(w: intellistream_datahub_sdk::resources::PolicyWarning) -> Self { + Self { + index: w.index, + external_id: w.external_id, + policy: w.policy, + message: w.message, + suggestion: w.suggestion, + } + } +} + +#[pymethods] +impl PyPolicyWarning { + fn __repr__(&self) -> String { + format!("PolicyWarning(external_id='{}')", self.external_id) + } +} diff --git a/datahub_python_bindings/src/resources/sync_service.rs b/datahub_python_bindings/src/resources/sync_service.rs index b4ad2e3..277ea5b 100644 --- a/datahub_python_bindings/src/resources/sync_service.rs +++ b/datahub_python_bindings/src/resources/sync_service.rs @@ -231,6 +231,54 @@ impl PyResourcesServiceSync { crate::PyPage::new(py, items, next_cursor) } + /// The whole connected graph component around one resource, as a gzip-compressed file + /// (`bytes`). It names everything by external id, so it imports into another tenant with + /// `import_graph`. Over 2,000,000 nodes or relationships raises a 400 rather than exporting + /// part. Use `export_graph_to_path` for anything large. + fn export_graph<'py>(&self, py: Python<'py>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + let file = py + .detach(|| self.runtime.block_on(service.resources.export_graph(id))) + .map_err(|e| crate::datahub_err(e))?; + Ok(pyo3::types::PyBytes::new(py, &file)) + } + + /// `export_graph`, streamed to `destination` without holding it in memory. Returns the number + /// of bytes written. + fn export_graph_to_path(&self, py: Python<'_>, id: u64, destination: String) -> PyResult { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.resources.export_graph_to_path(id, destination)) + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// Recreate the resources and relationships of an `export_graph` file. What already exists is + /// skipped, so re-importing into the source tenant is a no-op and a failed import can simply + /// be sent again. Timeseries are not created; missing ones are listed on the result. + fn import_graph(&self, py: Python<'_>, file: &[u8]) -> PyResult { + let service = self.api_service.clone(); + let file = file.to_vec(); + py.detach(|| { + self.runtime + .block_on(service.resources.import_graph(file)) + .map(Into::into) + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// `import_graph`, streaming the file from disk. + fn import_graph_from_path(&self, py: Python<'_>, source: String) -> PyResult { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.resources.import_graph_from_path(source)) + .map(Into::into) + .map_err(|e| crate::datahub_err(e)) + }) + } + /// `POST /resources/fetch-nearest` — the closest `limit` nodes carrying one of `end_labels`, /// plus the sub-graph connecting them back to the start. /// diff --git a/datahub_python_bindings/src/tenant/mod.rs b/datahub_python_bindings/src/tenant/mod.rs new file mode 100644 index 0000000..ed4bc95 --- /dev/null +++ b/datahub_python_bindings/src/tenant/mod.rs @@ -0,0 +1,315 @@ +use intellistream_datahub_sdk::tenant::{ + SettingsPermission, TenantFeatures, TenantLlmSettings, TenantLlmSettingsForm, +}; +use intellistream_datahub_sdk::ApiService; +use pyo3::prelude::*; +use pyo3_async_runtimes::tokio::future_into_py; +use std::collections::HashMap; +use std::sync::Arc; + +/// Which optional features are enabled for your tenant. +#[pyclass(module = "intellistream_datahub_sdk", name = "TenantFeatures", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyTenantFeatures { + pub files: bool, + pub policy: bool, + pub streaming: bool, + pub chat: bool, +} + +impl From for PyTenantFeatures { + fn from(f: TenantFeatures) -> Self { + Self { + files: f.files, + policy: f.policy, + streaming: f.streaming, + chat: f.chat, + } + } +} + +#[pymethods] +impl PyTenantFeatures { + fn __repr__(&self) -> String { + let b = |v: bool| if v { "True" } else { "False" }; + format!( + "TenantFeatures(files={}, policy={}, streaming={}, chat={})", + b(self.files), + b(self.policy), + b(self.streaming), + b(self.chat) + ) + } +} + +/// What the caller may do with one settings scope. +#[pyclass(module = "intellistream_datahub_sdk", name = "SettingsPermission", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PySettingsPermission { + pub read: bool, + pub write: bool, +} + +impl From for PySettingsPermission { + fn from(p: SettingsPermission) -> Self { + Self { + read: p.read, + write: p.write, + } + } +} + +#[pymethods] +impl PySettingsPermission { + fn __repr__(&self) -> String { + let b = |v: bool| if v { "True" } else { "False" }; + format!("SettingsPermission(read={}, write={})", b(self.read), b(self.write)) + } +} + +/// The model your organization's assistant runs on. The API key is never returned; `api_key_set` +/// says whether one is stored. +#[pyclass(module = "intellistream_datahub_sdk", name = "TenantLlmSettings", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyTenantLlmSettings { + pub provider: Option, + pub model: Option, + pub base_url: Option, + pub reasoning_effort: Option, + pub effort: Option, + pub turn_timeout: Option, + pub max_output_tokens: Option, + pub max_iterations: Option, + pub instructions: Option, + pub api_key_set: bool, + /// Whether this amounts to a model that can actually be called. + pub configured: bool, +} + +impl From for PyTenantLlmSettings { + fn from(s: TenantLlmSettings) -> Self { + Self { + provider: s.provider, + model: s.model, + base_url: s.base_url, + reasoning_effort: s.reasoning_effort, + effort: s.effort, + turn_timeout: s.turn_timeout, + max_output_tokens: s.max_output_tokens, + max_iterations: s.max_iterations, + instructions: s.instructions, + api_key_set: s.api_key_set, + configured: s.configured, + } + } +} + +#[pymethods] +impl PyTenantLlmSettings { + fn __repr__(&self) -> String { + format!( + "TenantLlmSettings(provider={:?}, model={:?}, configured={})", + self.provider, + self.model, + if self.configured { "True" } else { "False" } + ) + } +} + +fn permissions_to_py(map: HashMap) -> HashMap { + map.into_iter().map(|(k, v)| (k, v.into())).collect() +} + +#[allow(clippy::too_many_arguments)] +fn llm_form( + provider: Option, + model: Option, + api_key: Option, + base_url: Option, + reasoning_effort: Option, + effort: Option, + turn_timeout: Option, + max_output_tokens: Option, + max_iterations: Option, + instructions: Option, +) -> TenantLlmSettingsForm { + TenantLlmSettingsForm { + provider, + model, + api_key, + base_url, + reasoning_effort, + effort, + turn_timeout, + max_output_tokens, + max_iterations, + instructions, + } +} + +#[pyclass(module = "intellistream_datahub_sdk", name = "TenantServiceSync")] +pub struct PyTenantServiceSync { + pub api_service: Arc, + pub runtime: Arc, +} + +#[pymethods] +impl PyTenantServiceSync { + /// Which optional features are enabled for your tenant. + fn features(&self, py: Python<'_>) -> PyResult { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.tenant.features()) + .map(Into::into) + .map_err(crate::datahub_err) + }) + } + + /// What you may read and write, per settings scope (`"llm"`, …). For gating a UI; the settings + /// calls enforce the same grants. + fn settings_permissions(&self, py: Python<'_>) -> PyResult> { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.tenant.settings_permissions()) + .map(permissions_to_py) + .map_err(crate::datahub_err) + }) + } + + /// Your organization's model configuration. Needs the `llm` read grant. + fn llm_settings(&self, py: Python<'_>) -> PyResult { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.tenant.llm_settings()) + .map(Into::into) + .map_err(crate::datahub_err) + }) + } + + /// **Replace** the model configuration: an argument left out is cleared, except `api_key`, + /// where `None` or empty keeps the stored credential. Needs the `llm` write grant. + #[pyo3(signature = (provider=None, model=None, api_key=None, base_url=None, + reasoning_effort=None, effort=None, turn_timeout=None, + max_output_tokens=None, max_iterations=None, instructions=None))] + #[allow(clippy::too_many_arguments)] + fn update_llm_settings( + &self, + py: Python<'_>, + provider: Option, + model: Option, + api_key: Option, + base_url: Option, + reasoning_effort: Option, + effort: Option, + turn_timeout: Option, + max_output_tokens: Option, + max_iterations: Option, + instructions: Option, + ) -> PyResult { + let form = llm_form( + provider, model, api_key, base_url, reasoning_effort, effort, turn_timeout, + max_output_tokens, max_iterations, instructions, + ); + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.tenant.update_llm_settings(&form)) + .map(Into::into) + .map_err(crate::datahub_err) + }) + } +} + +#[pyclass(module = "intellistream_datahub_sdk", name = "TenantServiceAsync")] +pub struct PyTenantServiceAsync { + pub api_service: Arc, +} + +#[pymethods] +impl PyTenantServiceAsync { + /// Which optional features are enabled for your tenant. + fn features<'py>(&self, py: Python<'py>) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .tenant + .features() + .await + .map(PyTenantFeatures::from) + .map_err(crate::datahub_err) + }) + } + + /// What you may read and write, per settings scope (`"llm"`, …). + fn settings_permissions<'py>(&self, py: Python<'py>) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .tenant + .settings_permissions() + .await + .map(permissions_to_py) + .map_err(crate::datahub_err) + }) + } + + /// Your organization's model configuration. Needs the `llm` read grant. + fn llm_settings<'py>(&self, py: Python<'py>) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .tenant + .llm_settings() + .await + .map(PyTenantLlmSettings::from) + .map_err(crate::datahub_err) + }) + } + + /// **Replace** the model configuration: an argument left out is cleared, except `api_key`, + /// where `None` or empty keeps the stored credential. Needs the `llm` write grant. + #[pyo3(signature = (provider=None, model=None, api_key=None, base_url=None, + reasoning_effort=None, effort=None, turn_timeout=None, + max_output_tokens=None, max_iterations=None, instructions=None))] + #[allow(clippy::too_many_arguments)] + fn update_llm_settings<'py>( + &self, + py: Python<'py>, + provider: Option, + model: Option, + api_key: Option, + base_url: Option, + reasoning_effort: Option, + effort: Option, + turn_timeout: Option, + max_output_tokens: Option, + max_iterations: Option, + instructions: Option, + ) -> PyResult> { + let form = llm_form( + provider, model, api_key, base_url, reasoning_effort, effort, turn_timeout, + max_output_tokens, max_iterations, instructions, + ); + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .tenant + .update_llm_settings(&form) + .await + .map(PyTenantLlmSettings::from) + .map_err(crate::datahub_err) + }) + } +} + +pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> { + m.add_class::()?; + m.add_class::()?; + m.add_class::()?; + m.add_class::()?; + m.add_class::()?; + Ok(()) +} diff --git a/datahub_python_bindings/src/timeseries/async_service.rs b/datahub_python_bindings/src/timeseries/async_service.rs index 54f50fa..58b5fdb 100644 --- a/datahub_python_bindings/src/timeseries/async_service.rs +++ b/datahub_python_bindings/src/timeseries/async_service.rs @@ -45,6 +45,63 @@ impl PyTimeSeriesServiceAsync { }) } + /// One series by numeric id. Raises on 404 — and a 404 does not tell you the id is free: a + /// series you may not read is reported as missing rather than forbidden. + fn get_by_id<'p>(&self, py: Python<'p>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + let result = service + .time_series + .get_by_id(id) + .await + .map_err(|e| crate::datahub_err(e))?; + Ok(result + .get_items() + .first() + .map(|ts| PyTimeSeries::with_client(ts.clone(), service.clone()))) + }) + } + + /// The value type that compresses best for a unit while representing it faithfully. Advice, + /// not a constraint; an unknown unit answers the generic default with `recognized=False`. + fn recommend_value_type<'p>( + &self, + py: Python<'p>, + unit_external_id: String, + ) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + service + .time_series + .recommend_value_type(&unit_external_id) + .await + .map(crate::timeseries::live::PyValueTypeRecommendation::from) + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// Open a live tail of the datapoints written to these timeseries. Nothing is durable: the + /// stream starts at the latest point, and ids you cannot read are dropped silently. For + /// at-least-once delivery use `subscriptions.listen`. + #[pyo3(signature = (external_ids = Vec::new()))] + fn listen_datapoints<'p>( + &self, + py: Python<'p>, + external_ids: Vec, + ) -> PyResult> { + let service = self.api_service.clone(); + future_into_py(py, async move { + let listener = service + .time_series + .listen_datapoints(&external_ids) + .await + .map_err(crate::listen_err)?; + Ok(crate::timeseries::live::PyDatapointListenerAsync { + listener: crate::timeseries::live::shared_listener(listener), + }) + }) + } + fn create<'p>(&self, py: Python<'p>, input: Vec) -> PyResult> { let timeseries = input.iter().cloned().map(TimeSeries::from).collect(); let payload = DataWrapper::from_vec(timeseries); diff --git a/datahub_python_bindings/src/timeseries/live.rs b/datahub_python_bindings/src/timeseries/live.rs new file mode 100644 index 0000000..be5de8a --- /dev/null +++ b/datahub_python_bindings/src/timeseries/live.rs @@ -0,0 +1,277 @@ +use intellistream_datahub_sdk::timeseries::{DatapointListener, LiveDatapoint, ValueTypeRecommendation}; +use pyo3::exceptions::{PyStopAsyncIteration, PyStopIteration, PyValueError}; +use pyo3::prelude::*; +use pyo3_async_runtimes::tokio::future_into_py; +use std::sync::Arc; +use tokio::sync::Mutex; + +type SharedListener = Arc>>; + +pub(crate) fn shared_listener(l: DatapointListener) -> SharedListener { + Arc::new(Mutex::new(Some(l))) +} + +/// One point from `timeseries.listen_datapoints()`. `value` is a string for every value type. +#[pyclass(module = "intellistream_datahub_sdk", name = "LiveDatapoint", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyLiveDatapoint { + pub external_id: String, + pub value_type: Option, + pub timestamp: String, + pub value: String, +} + +impl From for PyLiveDatapoint { + fn from(p: LiveDatapoint) -> Self { + Self { + external_id: p.external_id, + value_type: p.value_type, + timestamp: p.timestamp, + value: p.value, + } + } +} + +#[pymethods] +impl PyLiveDatapoint { + fn __repr__(&self) -> String { + format!( + "LiveDatapoint(external_id='{}', timestamp='{}', value='{}')", + self.external_id, self.timestamp, self.value + ) + } +} + +/// Answer of `timeseries.recommend_value_type()`. +#[pyclass(module = "intellistream_datahub_sdk", name = "ValueTypeRecommendation", get_all, frozen, skip_from_py_object)] +#[derive(Clone)] +pub struct PyValueTypeRecommendation { + pub unit_external_id: String, + pub recommended_value_type: String, + pub reason: String, + pub recognized: bool, +} + +impl From for PyValueTypeRecommendation { + fn from(r: ValueTypeRecommendation) -> Self { + Self { + unit_external_id: r.unit_external_id, + recommended_value_type: r.recommended_value_type, + reason: r.reason, + recognized: r.recognized, + } + } +} + +#[pymethods] +impl PyValueTypeRecommendation { + fn __repr__(&self) -> String { + format!( + "ValueTypeRecommendation(unit_external_id='{}', recommended_value_type='{}', recognized={})", + self.unit_external_id, + self.recommended_value_type, + if self.recognized { "True" } else { "False" } + ) + } +} + +/// Synchronous live tail. `for point in listener:` blocks until the next datapoint. +#[pyclass(module = "intellistream_datahub_sdk", name = "DatapointListener")] +pub struct PyDatapointListener { + pub(crate) listener: SharedListener, + pub(crate) runtime: Arc, +} + +#[pymethods] +impl PyDatapointListener { + fn __iter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { + slf + } + + fn __next__(&self, py: Python<'_>) -> PyResult { + match self.next_datapoint(py)? { + Some(point) => Ok(point), + None => Err(PyStopIteration::new_err(())), + } + } + + /// Wait for the next datapoint; `None` once the listener is done. + fn next_datapoint(&self, py: Python<'_>) -> PyResult> { + let listener = self.listener.clone(); + let runtime = self.runtime.clone(); + py.detach(|| { + runtime.block_on(async move { + let mut guard = listener.lock().await; + let l = guard + .as_mut() + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; + match l.next().await { + Some(Ok(point)) => Ok(Some(PyLiveDatapoint::from(point))), + Some(Err(e)) => Err(crate::listen_err(e)), + None => Ok(None), + } + }) + }) + } + + /// Add timeseries to the live set. + fn subscribe(&self, py: Python<'_>, external_ids: Vec) -> PyResult<()> { + self.change(py, "subscribe", external_ids) + } + + /// Remove timeseries from the live set. + fn unsubscribe(&self, py: Python<'_>, external_ids: Vec) -> PyResult<()> { + self.change(py, "unsubscribe", external_ids) + } + + /// Replace the whole live set. + fn set_timeseries(&self, py: Python<'_>, external_ids: Vec) -> PyResult<()> { + self.change(py, "set", external_ids) + } + + fn close(&self, py: Python<'_>) -> PyResult<()> { + let listener = self.listener.clone(); + let runtime = self.runtime.clone(); + py.detach(|| { + runtime.block_on(async move { + if let Some(l) = listener.lock().await.take() { + l.close().await.map_err(crate::listen_err)?; + } + Ok(()) + }) + }) + } + + fn __enter__(slf: Py) -> Py { + slf + } + + #[pyo3(signature=(_exc_type=None, _exc_value=None, _traceback=None))] + fn __exit__<'py>( + &self, + py: Python<'py>, + _exc_type: Option>, + _exc_value: Option>, + _traceback: Option>, + ) -> PyResult<()> { + self.close(py) + } +} + +impl PyDatapointListener { + fn change(&self, py: Python<'_>, action: &'static str, ids: Vec) -> PyResult<()> { + let listener = self.listener.clone(); + let runtime = self.runtime.clone(); + py.detach(|| runtime.block_on(change_interest(listener, action, ids))) + } +} + +async fn change_interest( + listener: SharedListener, + action: &'static str, + ids: Vec, +) -> PyResult<()> { + let mut guard = listener.lock().await; + let l = guard + .as_mut() + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; + match action { + "subscribe" => l.subscribe(&ids).await, + "unsubscribe" => l.unsubscribe(&ids).await, + _ => l.set_timeseries(&ids).await, + } + .map_err(crate::listen_err) +} + +/// Asynchronous live tail. `async for point in listener:`. +#[pyclass(module = "intellistream_datahub_sdk", name = "DatapointListenerAsync")] +pub struct PyDatapointListenerAsync { + pub(crate) listener: SharedListener, +} + +#[pymethods] +impl PyDatapointListenerAsync { + fn __aiter__(slf: PyRef<'_, Self>) -> PyRef<'_, Self> { + slf + } + + fn __anext__<'py>(&self, py: Python<'py>) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + let mut guard = listener.lock().await; + let l = guard + .as_mut() + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; + match l.next().await { + Some(Ok(point)) => Ok(PyLiveDatapoint::from(point)), + Some(Err(e)) => Err(crate::listen_err(e)), + None => Err(PyStopAsyncIteration::new_err(())), + } + }) + } + + /// Wait for the next datapoint; `None` once the listener is done. + fn next_datapoint<'py>(&self, py: Python<'py>) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + let mut guard = listener.lock().await; + let l = guard + .as_mut() + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; + match l.next().await { + Some(Ok(point)) => Ok(Some(PyLiveDatapoint::from(point))), + Some(Err(e)) => Err(crate::listen_err(e)), + None => Ok(None), + } + }) + } + + fn subscribe<'py>(&self, py: Python<'py>, external_ids: Vec) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + change_interest(listener, "subscribe", external_ids).await?; + Ok(Python::attach(|py| py.None())) + }) + } + + fn unsubscribe<'py>(&self, py: Python<'py>, external_ids: Vec) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + change_interest(listener, "unsubscribe", external_ids).await?; + Ok(Python::attach(|py| py.None())) + }) + } + + fn set_timeseries<'py>(&self, py: Python<'py>, external_ids: Vec) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + change_interest(listener, "set", external_ids).await?; + Ok(Python::attach(|py| py.None())) + }) + } + + fn close<'py>(&self, py: Python<'py>) -> PyResult> { + let listener = self.listener.clone(); + future_into_py(py, async move { + if let Some(l) = listener.lock().await.take() { + l.close().await.map_err(crate::listen_err)?; + } + Ok(Python::attach(|py| py.None())) + }) + } + + fn __aenter__<'py>(slf: Py, py: Python<'py>) -> PyResult> { + future_into_py(py, async move { Ok(slf) }) + } + + #[pyo3(signature=(_exc_type=None, _exc_value=None, _traceback=None))] + fn __aexit__<'py>( + &self, + py: Python<'py>, + _exc_type: Option>, + _exc_value: Option>, + _traceback: Option>, + ) -> PyResult> { + self.close(py) + } +} diff --git a/datahub_python_bindings/src/timeseries/mod.rs b/datahub_python_bindings/src/timeseries/mod.rs index cdd8424..5cd6d94 100644 --- a/datahub_python_bindings/src/timeseries/mod.rs +++ b/datahub_python_bindings/src/timeseries/mod.rs @@ -36,6 +36,7 @@ pub mod async_service; mod construction; pub mod datapoints; pub mod general; +pub mod live; pub mod sync_service; /// Python wrapper for Timeseries objects, represents contextualization data for timeseries @@ -469,6 +470,10 @@ pub fn register(m: &Bound<'_, PyModule>) -> PyResult<()> { m.add_class::()?; m.add_class::()?; m.add_class::()?; + m.add_class::()?; + m.add_class::()?; + m.add_class::()?; + m.add_class::()?; Ok(()) } diff --git a/datahub_python_bindings/src/timeseries/sync_service.rs b/datahub_python_bindings/src/timeseries/sync_service.rs index 325aec5..0d84537 100644 --- a/datahub_python_bindings/src/timeseries/sync_service.rs +++ b/datahub_python_bindings/src/timeseries/sync_service.rs @@ -41,6 +41,57 @@ impl PyTimeSeriesServiceSync { }) } + /// One series by numeric id. Raises on 404 — and a 404 does not tell you the id is free: a + /// series you may not read is reported as missing rather than forbidden. + fn get_by_id(&self, py: Python<'_>, id: u64) -> PyResult> { + let service = self.api_service.clone(); + let result = py.detach(|| self.runtime.block_on(service.time_series.get_by_id(id))); + let result = result.map_err(|e| crate::datahub_err(e))?; + Ok(result + .get_items() + .first() + .map(|ts| PyTimeSeries::with_client(ts.clone(), service.clone()))) + } + + /// The value type that compresses best for a unit while representing it faithfully. Advice, + /// not a constraint; an unknown unit answers the generic default with `recognized=False`. + fn recommend_value_type( + &self, + py: Python<'_>, + unit_external_id: String, + ) -> PyResult { + let service = self.api_service.clone(); + py.detach(|| { + self.runtime + .block_on(service.time_series.recommend_value_type(&unit_external_id)) + .map(Into::into) + .map_err(|e| crate::datahub_err(e)) + }) + } + + /// Open a live tail of the datapoints written to these timeseries. Nothing is durable: the + /// stream starts at the latest point, and ids you cannot read are dropped silently. For + /// at-least-once delivery use `subscriptions.listen`. + #[pyo3(signature = (external_ids = Vec::new()))] + fn listen_datapoints( + &self, + py: Python<'_>, + external_ids: Vec, + ) -> PyResult { + let service = self.api_service.clone(); + let runtime = self.runtime.clone(); + py.detach(|| { + let listener = self + .runtime + .block_on(service.time_series.listen_datapoints(&external_ids)) + .map_err(crate::listen_err)?; + Ok(crate::timeseries::live::PyDatapointListener { + listener: crate::timeseries::live::shared_listener(listener), + runtime, + }) + }) + } + fn create<'p>(&self, py: Python<'p>, input: Vec) -> PyResult> { let timeseries = input.iter().cloned().map(TimeSeries::from).collect(); let payload = DataWrapper::from_vec(timeseries); diff --git a/python_tests/test_functions.py b/python_tests/test_functions.py index d5c630c..b4a1dba 100644 --- a/python_tests/test_functions.py +++ b/python_tests/test_functions.py @@ -6,7 +6,7 @@ import intellistream_datahub_sdk import pytest -from fixtures import sync_client, unique_id +from fixtures import make_function, sync_client, unique_id def test_create_list_by_external_id_delete(sync_client): @@ -45,3 +45,32 @@ def test_create_list_by_external_id_delete(sync_client): def test_by_external_id_raises_when_missing(sync_client): with pytest.raises(Exception): sync_client.functions.by_external_id(unique_id("does_not_exist")) + + +def test_by_ids_filter_and_search(sync_client, make_function): + from polling import poll_until + + # One unbroken lexeme for the full-text index; underscores would split it. + token = unique_id("fn").replace("_", "") + fn = make_function(name=f"{token} filter probe") + + by_id = sync_client.functions.by_ids([fn.id]) + assert [f.external_id for f in by_id] == [fn.external_id] + assert sync_client.functions.by_ids([unique_id("fn_absent")]) == [] + + page = sync_client.functions.filter(external_id=fn.external_id) + assert [f.external_id for f in page] == [fn.external_id] + same = sync_client.functions.filter( + filter=intellistream_datahub_sdk.FunctionFilter(external_id=fn.external_id) + ) + assert [f.external_id for f in same] == [fn.external_id] + with pytest.raises(TypeError): + sync_client.functions.filter( + filter=intellistream_datahub_sdk.FunctionFilter(), external_id=fn.external_id + ) + + hits = poll_until( + lambda: sync_client.functions.search(token), + lambda found: any(f.external_id == fn.external_id for f in found), + ) + assert any(f.external_id == fn.external_id for f in hits) diff --git a/python_tests/test_graph_transfer.py b/python_tests/test_graph_transfer.py new file mode 100644 index 0000000..be00d11 --- /dev/null +++ b/python_tests/test_graph_transfer.py @@ -0,0 +1,64 @@ +"""Tests for `resources.export_graph` / `import_graph`. Mirrors the graph-transfer tests in +`src/resources/tests.rs`.""" +import intellistream_datahub_sdk +import pytest + +from fixtures import async_client, make_resource, sync_client, unique_id +from polling import poll_until + + +def _component(make_resource, sync_client): + root, child = unique_id("export_root"), unique_id("export_child") + created = make_resource( + [ + intellistream_datahub_sdk.Resource( + external_id=root, name="python export root", labels=["ASSET"], is_root=True + ), + intellistream_datahub_sdk.Resource( + external_id=child, name="python export child", labels=["ASSET"] + ), + ], + [ + intellistream_datahub_sdk.RelForm( + relationship_type="FLOWS_TO", from_external_id=root, to_external_id=child + ) + ], + ) + # The export walks the graph projection, which lags the write. + network = poll_until( + lambda: sync_client.resources.fetch_related(root, depth=-1), + lambda n: len(n.nodes) >= 2, + ) + assert len(network.nodes) >= 2, "graph projection did not catch up" + return next(n.id for n in created.nodes if n.external_id == root) + + +def test_export_then_import_into_the_source_tenant_is_a_no_op(sync_client, make_resource, tmp_path): + root_id = _component(make_resource, sync_client) + + file = sync_client.resources.export_graph(root_id) + assert isinstance(file, bytes) + assert file[:2] == b"\x1f\x8b", "the export is gzip" + + result = sync_client.resources.import_graph(file) + assert result.nodes_created == 0 + assert result.nodes_skipped_existing >= 2 + + path = tmp_path / "component.graph" + written = sync_client.resources.export_graph_to_path(root_id, str(path)) + assert written == path.stat().st_size > 0 + from_disk = sync_client.resources.import_graph_from_path(str(path)) + assert from_disk.nodes_created == 0 + + +def test_export_of_an_unknown_id_raises_404(sync_client): + with pytest.raises(intellistream_datahub_sdk.DataHubException) as err: + sync_client.resources.export_graph(2**62) + assert err.value.status_code == 404 + + +@pytest.mark.asyncio +async def test_async_export_returns_bytes(async_client, sync_client, make_resource): + root_id = _component(make_resource, sync_client) + file = await async_client.resources.export_graph(root_id) + assert isinstance(file, bytes) and file[:2] == b"\x1f\x8b" diff --git a/python_tests/test_resources.py b/python_tests/test_resources.py index 8d509c1..a57171a 100644 --- a/python_tests/test_resources.py +++ b/python_tests/test_resources.py @@ -10,7 +10,7 @@ import pytest from intellistream_datahub_sdk import DataHubException, EdgeProxy, GraphResult, RelForm, Resource -from fixtures import sync_client, unique_id +from fixtures import TEST_LABEL, make_resource, sync_client, unique_id from polling import poll_until @@ -161,3 +161,39 @@ def test_api_error_surfaces_status_code(sync_client): sync_client.resources.create([bad]) assert exc_info.value.status_code == 400 assert exc_info.value.message # raw response body is preserved + + +def test_fetch_nearest_reaches_the_labelled_node_through_the_path(sync_client, make_resource): + """Mirrors `fetch_nearest_reaches_the_labelled_node_through_the_path` in + `src/resources/tests.rs`: `root -> middle -> leaf`, only the leaf labelled.""" + root, middle, leaf = (unique_id(k) for k in ("nearest_root", "nearest_middle", "nearest_leaf")) + created = make_resource( + [ + Resource(external_id=root, name="py nearest root", labels=["ASSET"], is_root=True), + Resource(external_id=middle, name="py nearest middle", labels=["ASSET"]), + Resource(external_id=leaf, name="py nearest leaf", labels=["ASSET", TEST_LABEL]), + ], + [ + RelForm(relationship_type="flows_to", from_external_id=root, to_external_id=middle), + RelForm(relationship_type="flows_to", from_external_id=middle, to_external_id=leaf), + ], + ) + root_id = next(n.id for n in created.nodes if n.external_id == root) + + def reached(network, ext): + return any(n.external_id == ext for n in network.nodes) + + network = poll_until( + lambda: sync_client.resources.fetch_nearest(root_id, end_labels=[TEST_LABEL], limit=1), + lambda n: reached(n, leaf), + ) + assert reached(network, leaf), "the labelled leaf was not reached" + assert reached(network, middle), "the path back to the start is part of the answer" + + none = sync_client.resources.fetch_nearest( + root_id, + end_labels=[TEST_LABEL], + limit=1, + relationship_types=[unique_id("no_such_type").upper()], + ) + assert not reached(none, leaf) diff --git a/python_tests/test_single_reads.py b/python_tests/test_single_reads.py new file mode 100644 index 0000000..08968ab --- /dev/null +++ b/python_tests/test_single_reads.py @@ -0,0 +1,65 @@ +"""By-id reads, the value-type recommendation, and the live datapoint tail.""" +import threading + +import intellistream_datahub_sdk +import pandas as pd +import pytest + +from fixtures import async_client, make_dataset, make_ts, sync_client, unique_id + + +def test_timeseries_get_by_id(sync_client, make_ts): + ts = make_ts() + fetched = sync_client.timeseries.get_by_id(ts.id) + assert fetched.external_id == ts.external_id + + +def test_timeseries_get_by_id_of_an_unknown_id_raises_404(sync_client): + with pytest.raises(intellistream_datahub_sdk.DataHubException) as err: + sync_client.timeseries.get_by_id(2**62) + assert err.value.status_code == 404 + assert err.value.problem_slug == "not-found" + + +def test_dataset_get_by_id(sync_client, make_dataset): + ds = make_dataset() + fetched = sync_client.datasets.get_by_id(ds.id) + assert fetched.external_id == ds.external_id + + +def test_recommend_value_type(sync_client): + known = sync_client.timeseries.recommend_value_type("temperature_deg_c") + assert known.unit_external_id == "temperature_deg_c" + assert known.recognized + unknown = sync_client.timeseries.recommend_value_type(unique_id("unit")) + assert not unknown.recognized + + +@pytest.mark.asyncio +async def test_async_get_by_id(async_client, make_ts): + ts = make_ts() + fetched = await async_client.timeseries.get_by_id(ts.id) + assert fetched.external_id == ts.external_id + + +def test_listen_datapoints_delivers_points_written_after_connecting(sync_client, make_ts): + ts = make_ts() + received = [] + listener = sync_client.timeseries.listen_datapoints([ts.external_id]) + + reader = threading.Thread(target=lambda: received.append(listener.next_datapoint()), daemon=True) + reader.start() + # The server-side consumer reads from latest, so write after it has attached. + reader.join(timeout=2) + sync_client.timeseries.insert_from_lists( + timestamps=pd.DatetimeIndex([pd.Timestamp.now(tz="UTC")]), values=[21.5], ts=ts + ) + reader.join(timeout=30) + # A reader still waiting holds the listener, so closing now would block on it rather than + # fail; leave the daemon thread and its socket to the interpreter. + if not reader.is_alive(): + listener.close() + + assert received, "no datapoint within 30s" + assert received[0].external_id == ts.external_id + assert float(received[0].value) == 21.5 diff --git a/python_tests/test_tenant.py b/python_tests/test_tenant.py new file mode 100644 index 0000000..a8bb306 --- /dev/null +++ b/python_tests/test_tenant.py @@ -0,0 +1,64 @@ +"""Tests for the `/tenant` bindings.""" +import intellistream_datahub_sdk +import pytest + +from fixtures import async_client, sync_client + + +def test_features_answer_every_flag(sync_client): + features = sync_client.tenant.features() + for flag in ("files", "policy", "streaming", "chat"): + assert isinstance(getattr(features, flag), bool) + + +def test_llm_settings_follow_the_reported_permission(sync_client): + permissions = sync_client.tenant.settings_permissions() + assert "llm" in permissions + if permissions["llm"].read: + settings = sync_client.tenant.llm_settings() + assert isinstance(settings.configured, bool) + else: + with pytest.raises(intellistream_datahub_sdk.DataHubException) as err: + sync_client.tenant.llm_settings() + assert err.value.status_code == 403 + + +@pytest.mark.asyncio +async def test_async_features(async_client): + features = await async_client.tenant.features() + assert isinstance(features.files, bool) + + +def test_update_llm_settings_is_gated_validated_and_round_trips(sync_client): + """Mirrors the Rust test: 403 without the write grant; with it, an empty form is a 400 naming + `provider` and `model`, and a configured model written back unchanged answers what was stored. + The grant is checked before the form, so no branch changes the tenant's settings.""" + permissions = sync_client.tenant.settings_permissions() + if not permissions["llm"].write: + with pytest.raises(intellistream_datahub_sdk.DataHubException) as err: + sync_client.tenant.update_llm_settings() + assert err.value.status_code == 403 + return + + with pytest.raises(intellistream_datahub_sdk.DataHubException) as err: + sync_client.tenant.update_llm_settings() + assert err.value.status_code == 400 + fields = {f.get("field") for f in (err.value.problem or {}).get("fields", [])} + assert {"provider", "model"} <= fields + + stored = sync_client.tenant.llm_settings() + if not stored.configured: + pytest.skip("no model configured; nothing to write back unchanged") + saved = sync_client.tenant.update_llm_settings( + provider=stored.provider, + model=stored.model, + base_url=stored.base_url, + reasoning_effort=stored.reasoning_effort, + effort=stored.effort, + turn_timeout=stored.turn_timeout, + max_output_tokens=stored.max_output_tokens, + max_iterations=stored.max_iterations, + instructions=stored.instructions, + ) + for attr in ("provider", "model", "base_url", "effort", "instructions", "api_key_set"): + assert getattr(saved, attr) == getattr(stored, attr), attr diff --git a/src/blocking.rs b/src/blocking.rs index aec476c..1dd827a 100644 --- a/src/blocking.rs +++ b/src/blocking.rs @@ -32,7 +32,9 @@ use crate::datasets::{DatasetFilter, Dataset, DatasetFilterForm, DatasetUpdate}; use crate::events::{Event, EventDimension, EventIdCollection}; use crate::files::{FileDownload, FileUpdate, FileUpload}; use crate::filters::{EventFilter, EventFilterForm}; -use crate::functions::Function; +use crate::functions::{Function, FunctionFilter, FunctionFilterForm}; +use crate::tenant::{SettingsPermission, TenantFeatures, TenantLlmSettings, TenantLlmSettingsForm}; +use std::collections::HashMap; use crate::generic::{ DataWrapper, Datapoint, DatapointString, DatapointsCollection, DeleteFilter, INode, IdAndExtId, RetrieveFilter, SearchAndFilterForm, @@ -43,10 +45,13 @@ use crate::labels::Label; use crate::nodes::{Asset, Node}; use crate::relations::{EdgeProxy, RelForm, RelTypeForm, RelationshipType}; use crate::resources::{ - RelatedResourcesForm, Resource, ResourceFilter, ResourceFilterForm, ResourceNetwork, - ResourceUpdate, + GraphImportResult, RelatedResourcesForm, Resource, ResourceFilter, ResourceFilterForm, + ResourceNetwork, ResourceUpdate, +}; +use crate::timeseries::{ + BinaryIngestOptions, TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection, + ValueTypeRecommendation, }; -use crate::timeseries::{BinaryIngestOptions, TimeSeries, TimeSeriesFilter, TimeSeriesUpdateCollection}; use crate::unit::Unit; /// Generate blocking methods that delegate to the same-named async method on one of @@ -100,6 +105,7 @@ pub struct ApiService { pub functions: FunctionsService, pub labels: LabelsService, pub edges: EdgesService, + pub tenant: TenantService, } /// The blocking counterpart of [`crate::create_api_service`]: configuration from the @@ -141,6 +147,7 @@ impl ApiService { functions: service!(FunctionsService), labels: service!(LabelsService), edges: service!(EdgesService), + tenant: service!(TenantService), api, } } @@ -161,6 +168,8 @@ pub struct TimeSeriesService { impl TimeSeriesService { delegate! { time_series => fn list(limit: Option) -> Result, ResponseError>; + fn get_by_id(id: u64) -> Result, ResponseError>; + fn recommend_value_type(unit_external_id: &str) -> Result; fn create(json: &DataWrapper) -> Result, ResponseError>; fn create_one(ts: &TimeSeries) -> Result, ResponseError>; fn create_from_list(ts_list: &Vec) -> Result, ResponseError>; @@ -199,6 +208,27 @@ impl ResourceService { fn list(limit: Option) -> Result, ResponseError>; fn search(payload: &SearchAndFilterForm) -> Result, ResponseError>; fn fetch_related(form: &RelatedResourcesForm) -> Result; + fn export_graph(id: u64) -> Result, ResponseError>; + fn import_graph(file: Vec) -> Result; + } + + /// Blocking counterpart of [`crate::ResourceService::export_graph_to_path`]. + pub fn export_graph_to_path( + &self, + id: u64, + destination: impl AsRef, + ) -> Result { + self.rt + .block_on(self.api.resources.export_graph_to_path(id, destination)) + } + + /// Blocking counterpart of [`crate::ResourceService::import_graph_from_path`]. + pub fn import_graph_from_path( + &self, + source: impl AsRef, + ) -> Result { + self.rt + .block_on(self.api.resources.import_graph_from_path(source)) } // Generic or GraphDataWrapper-returning; delegated by hand. @@ -277,6 +307,7 @@ pub struct DatasetsService { impl DatasetsService { delegate! { datasets => fn list(limit: Option) -> Result, ResponseError>; + fn get_by_id(id: u64) -> Result, ResponseError>; fn filter(filter: &DatasetFilterForm) -> Result, ResponseError>; fn search(search: &SearchAndFilterForm) -> Result, ResponseError>; fn search_by_query(query: &str) -> Result, ResponseError>; @@ -385,6 +416,9 @@ impl FunctionsService { fn list(limit: Option) -> Result, ResponseError>; fn by_ids(ids: &[IdAndExtId]) -> Result, ResponseError>; fn by_external_id(external_id: &str) -> Result; + fn filter(form: &FunctionFilterForm) -> Result, ResponseError>; + fn search(form: &SearchAndFilterForm) -> Result, ResponseError>; + fn search_by_query(query: &str) -> Result, ResponseError>; } delegate_into! { functions => @@ -439,3 +473,18 @@ impl EdgesService { fn create_types(data: Into>) -> Result, ResponseError>; } } + +/// Blocking counterpart of [`crate::tenant::TenantService`]. +pub struct TenantService { + api: Arc, + rt: Arc, +} + +impl TenantService { + delegate! { tenant => + fn features() -> Result; + fn settings_permissions() -> Result, ResponseError>; + fn llm_settings() -> Result; + fn update_llm_settings(form: &TenantLlmSettingsForm) -> Result; + } +} diff --git a/src/datasets/mod.rs b/src/datasets/mod.rs index eacdbf9..2dba8e0 100644 --- a/src/datasets/mod.rs +++ b/src/datasets/mod.rs @@ -70,6 +70,14 @@ impl DatasetsService { .await } + /// `GET /datasets/{id}` — one data set by its numeric id. A miss is a 404 carrying the + /// `not-found` problem type, unlike [`by_ids`](Self::by_ids), which omits what it cannot find. + pub async fn get_by_id(&self, id: u64) -> Result, ResponseError> { + let path = &format!("{}/{}", self.base_url, id); + self.execute_get_request::, ()>(path, None) + .await + } + /// `POST /datasets/filter` — datasets matching [`DatasetFilterForm`], newest first. /// /// Every criterion on the filter is honoured server-side. Results are capped by the form's diff --git a/src/datasets/tests.rs b/src/datasets/tests.rs index 02c17fb..6992652 100644 --- a/src/datasets/tests.rs +++ b/src/datasets/tests.rs @@ -307,3 +307,25 @@ fn filter_body_matches_the_documented_wire_shape() { "ids must be strings so a large id survives a JavaScript client" ); } + +#[tokio::test] +async fn get_by_id_reads_one_dataset_and_404s_a_miss() -> Result<(), ResponseError> { + let api = create_api_service(); + let dataset = Dataset::new(unique_id("dataset")).build(); + let ext_id = dataset.external_id().to_string(); + let created = api.datasets.create(&vec![dataset]).await?; + let _cleanup = cleanup_datasets(vec![ext_id.clone()]); + let id = *created.get_items()[0].id().expect("create echoes the id"); + + let got = api.datasets.get_by_id(id).await?; + assert_eq!(got.get_items().len(), 1); + assert_eq!(got.get_items()[0].external_id(), &ext_id); + + let err = api + .datasets + .get_by_id(u64::MAX / 2) + .await + .expect_err("an unknown id is a 404"); + assert_eq!(err.status.as_u16(), 404); + Ok(()) +} diff --git a/src/functions/mod.rs b/src/functions/mod.rs index fb51a85..557a246 100644 --- a/src/functions/mod.rs +++ b/src/functions/mod.rs @@ -1,7 +1,8 @@ #[cfg(test)] mod test; -use crate::generic::{ApiServiceProvider, DataHubEntity, DataWrapper, IdAndExtId}; +use crate::filters::{NodeFilter, PageRequest}; +use crate::generic::{ApiServiceProvider, DataHubEntity, DataWrapper, IdAndExtId, SearchAndFilterForm}; use crate::graph_data_wrapper::GraphDataWrapper; use crate::http::ResponseError; use crate::nodes::Node; @@ -66,9 +67,7 @@ impl FunctionsService { /// missing rather than forbidden, so a 404 says "not a function you can read" and nothing /// more. A node of another type is not a function and is reported the same way. /// - /// Unlike [`by_ids`](Self::by_ids), which omits what it cannot find, this is an error. Prefer - /// it to `by_ids` when you already have the numeric id: `by_ids` does not yet call the api's - /// `/functions/byids` and pages the whole listing to filter client-side. + /// Unlike [`by_ids`](Self::by_ids), which omits what it cannot find, this is an error. pub async fn get_by_id(&self, id: u64) -> Result, ResponseError> { let path = &format!("{}/{}", self.base_url, id); self.execute_get_request::, ()>(path, None) @@ -105,43 +104,58 @@ impl FunctionsService { .await } - /// Look up functions by id or externalId, implemented client-side by listing and filtering. + /// `POST /functions/byids` — a batch lookup by id or external id. /// - /// **The SDK has not wired the real endpoint yet.** The api grew `/functions/byids`, - /// `/functions/filter` and `/functions/search` in platform #131, which is what this should - /// call; until it does, the client-side walk stands. - /// - /// It asks for the largest page the api allows, because a client-side filter can only match - /// what the listing returned — so a tenant past 10000 functions silently misses the oldest - /// ones here. The fix is to call `/functions/byids`, not to ask for a bigger page. + /// Like every batch lookup in this api, it answers 200 with the found subset and silently + /// omits the rest: an id that does not exist, names a node of another type, or is not readable + /// is simply absent from the response. pub async fn by_ids( &self, ids: &[IdAndExtId], ) -> Result, ResponseError> { - let mut wanted_ids: Vec = vec![]; - let mut wanted_external_ids: Vec = vec![]; - for id in ids { - if let Some(numeric) = id.id { - wanted_ids.push(numeric); - } - if let Some(ext) = &id.external_id { - wanted_external_ids.push(ext.clone()); - } - } - let all = self.list(Some(10_000)).await?; - let mut matched: Vec = vec![]; - for f in all.get_items() { - let id_match = f.id.map_or(false, |i| wanted_ids.contains(&i)); - let ext_match = wanted_external_ids.contains(&f.external_id); - if id_match || ext_match { - matched.push(f.clone()); - } - } - let mut wrapper = DataWrapper::from_vec(matched); - if let Some(code) = all.get_http_status_code() { - wrapper.set_http_status_code(code); - } - Ok(wrapper) + let path = &format!("{}/byids", self.base_url); + let body: DataWrapper = DataWrapper::from_vec(ids.to_vec()); + self.execute_post_request::, _>(path, &body) + .await + } + + /// `POST /functions/filter` — functions matching every supplied criterion, newest created + /// first unless the form sorts otherwise. + /// + /// The shared [`NodeFilter`](crate::filters::NodeFilter) rules apply — wildcards, + /// case-insensitivity, what an empty list means — plus + /// [`data_set_id`](FunctionFilter::data_set_id). Pages: the form carries `sort` and `cursor`, + /// and the response carries `next_cursor`. + pub async fn filter( + &self, + form: &FunctionFilterForm, + ) -> Result, ResponseError> { + let path = &format!("{}/filter", self.base_url); + self.execute_post_request::, _>(path, form) + .await + } + + /// `POST /functions/search` — free-text search over functions, ranked by `ts_rank` and + /// tie-broken on id. + /// + /// The phrase selects and the filter only removes. `query` is required at 3–140 characters; + /// `limit` defaults to 100 and caps at 1000. + pub async fn search( + &self, + form: &SearchAndFilterForm, + ) -> Result, ResponseError> { + let path = &format!("{}/search", self.base_url); + self.execute_post_request::, _>(path, form) + .await + } + + /// [`search`](Self::search) with just a query string, leaving `limit` at the server's default + /// of 100. + pub async fn search_by_query( + &self, + query: &str, + ) -> Result, ResponseError> { + self.search(&SearchAndFilterForm::new(query)).await } /// Convenience for the function-worker bootstrap: `client.functions.by_external_id("...")`. @@ -253,3 +267,52 @@ impl DataHubEntity for Function { &self.external_id } } + +/// Criteria for `POST /functions/filter`, and the `filter` of `POST /functions/search`: the shared +/// [`NodeFilter`] plus a data set restriction, mirroring the api's `FunctionFilter`. +// Not PartialEq: `data_set_id` holds `IdAndExtId`, which is intentionally non-comparable. +#[derive(Debug, Serialize, Deserialize, Clone, Default)] +#[serde(rename_all = "camelCase")] +pub struct FunctionFilter { + #[serde(flatten)] + pub node: NodeFilter, + /// Restrict to functions in these data sets **and every data set beneath them**, each named + /// by id or external id. + /// + /// **`None` and empty differ**: `None` places no restriction, `Some(vec![])` narrows to no + /// data sets and matches nothing. + #[serde(skip_serializing_if = "Option::is_none")] + pub data_set_id: Option>, +} + +/// Body of `POST /functions/filter`: the criteria, how many to return, and in what order. +#[derive(Debug, Serialize, Deserialize, Clone, Default)] +#[serde(rename_all = "camelCase")] +pub struct FunctionFilterForm { + pub filter: FunctionFilter, + /// Defaults to 1000 server-side and is capped at 10000 — above that the request is a 400. + #[serde(skip_serializing_if = "Option::is_none")] + pub limit: Option, + #[serde(flatten)] + pub paging: PageRequest, +} + +impl FunctionFilterForm { + pub fn new(filter: FunctionFilter) -> Self { + Self { + filter, + limit: None, + paging: Default::default(), + } + } + + pub fn with_limit(mut self, limit: u64) -> Self { + self.limit = Some(limit); + self + } + + pub fn with_paging(mut self, paging: PageRequest) -> Self { + self.paging = paging; + self + } +} diff --git a/src/functions/test.rs b/src/functions/test.rs index e1ead80..1c8e42c 100644 --- a/src/functions/test.rs +++ b/src/functions/test.rs @@ -94,4 +94,56 @@ mod tests { let err = after.expect_err("a deleted function is a 404"); assert_eq!(err.status.as_u16(), 404); } + + /// `/functions/byids`, `/filter` and `/search` — the server-side reads that replaced the + /// client-side walk over the listing. + #[tokio::test] + #[ignore] + async fn functions_byids_filter_and_search() { + use crate::filters::NodeFilter; + use crate::functions::{FunctionFilter, FunctionFilterForm}; + use crate::tests::ids::unique_token; + use crate::tests::polling::poll_until; + + let api = create_api_service(); + let ext_id = unique_id("fn"); + let token = unique_token("fn"); + let mut function_in = Function::new(ext_id.clone()).with_name(format!("{token} filter probe")); + function_in.description = Some(format!("{token} searchable description")); + + let created = api.functions.create(&vec![function_in]).await.unwrap(); + let _cleanup = cleanup_functions(vec![ext_id.clone()]); + let id = created.get_items()[0].id.unwrap(); + + let by_id = api.functions.by_ids(&[IdAndExtId::from_id(id)]).await.unwrap(); + assert_eq!(by_id.get_items().len(), 1); + assert_eq!(by_id.get_items()[0].external_id, ext_id); + let missing = api + .functions + .by_ids(&[IdAndExtId::from_external_id(&unique_id("fn_absent"))]) + .await + .unwrap(); + assert!(missing.get_items().is_empty(), "a batch lookup omits what it cannot find"); + + let form = FunctionFilterForm::new(FunctionFilter { + node: NodeFilter { + external_id: Some(vec![ext_id.clone()]), + ..Default::default() + }, + data_set_id: None, + }); + let filtered = api.functions.filter(&form).await.unwrap(); + assert_eq!(filtered.get_items().len(), 1); + assert_eq!(filtered.get_items()[0].external_id, ext_id); + + let found = poll_until( + || async { api.functions.search_by_query(&token).await.unwrap() }, + |hits| hits.get_items().iter().any(|f| f.external_id == ext_id), + ) + .await; + assert!( + found.get_items().iter().any(|f| f.external_id == ext_id), + "search for {token} never returned {ext_id}" + ); + } } diff --git a/src/generic.rs b/src/generic.rs index feddc72..88c7354 100644 --- a/src/generic.rs +++ b/src/generic.rs @@ -804,6 +804,32 @@ pub trait ApiServiceProvider { } } + /// `PUT` a JSON body. Only the tenant settings replace a whole object this way; every other + /// write in this api is a `POST`. + async fn execute_put_request( + &self, + path: &str, + json: &J, + ) -> Result { + let token = self.get_token().await?; + let response = self + .get_api_service() + .http_client + .put(path) + .json(json) + .bearer_auth(token.clone()) + .send() + .await + .map_err(|err| { + eprintln!("HTTP request failed: {}", err); + ResponseError::from_err(err) + })?; + match process_response::(response, path).await { + Ok(value) => Ok(value), + Err(e) => Err(self.on_request_error(e, &token).await), + } + } + /// Uploads a file with a raw `PUT`: the file content is the request body and all metadata /// travels in headers (`X-Datahub-Path`, `X-Datahub-External-Id`, `X-Datahub-Dataset-Id`, /// `X-Datahub-Description`, `Content-Type`). The server validates and authorises the upload @@ -837,12 +863,12 @@ pub trait ApiServiceProvider { } /// `POST` a binary body under its own media type: the datapoint frames of - /// `/timeseries/data/binary`. The api's 204 becomes an empty wrapper, as in the JSON helper; - /// a rejection carries the api's `problem+json` text. + /// `/timeseries/data/binary` and the graph files of `/resources/import`. The api's 204 becomes + /// an empty wrapper, as in the JSON helper; a rejection carries the api's `problem+json` text. async fn execute_post_bytes_request( &self, path: &str, - body: Vec, + body: impl Into, content_type: &str, ) -> Result { let token = self.get_token().await?; diff --git a/src/lib.rs b/src/lib.rs index 23fd824..44f4915 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -52,6 +52,7 @@ pub mod tests; pub mod timeseries; pub mod unit; pub mod functions; +pub mod tenant; pub use resources::*; pub use nodes::{Asset, Node, NodeType, Policy}; @@ -72,6 +73,7 @@ pub use subscriptions::{ SubscriptionFilterForm, WsDatapoint, }; use crate::functions::FunctionsService; +pub use crate::tenant::TenantService; //pub use filters::Filter; pub struct ApiService { @@ -87,6 +89,7 @@ pub struct ApiService { pub functions: FunctionsService, pub labels: LabelsService, pub edges: EdgesService, + pub tenant: TenantService, pub(crate) http_client: Client, } @@ -144,6 +147,7 @@ pub fn create_api_service() -> Arc { functions: FunctionsService::new(Weak::clone(weak_self), &base_url_clone), labels: LabelsService::new(Weak::clone(weak_self), &base_url_clone), edges: EdgesService::new(Weak::clone(weak_self), &base_url_clone), + tenant: TenantService::new(Weak::clone(weak_self), &base_url_clone), http_client, } }); @@ -181,6 +185,7 @@ impl ApiService { functions: FunctionsService::new(Weak::clone(weak_self), &base_url_clone), labels: LabelsService::new(Weak::clone(weak_self), &base_url_clone), edges: EdgesService::new(Weak::clone(weak_self), &base_url_clone), + tenant: TenantService::new(Weak::clone(weak_self), &base_url_clone), http_client, } }); @@ -219,6 +224,7 @@ impl ApiService { functions: FunctionsService::new(Weak::clone(weak_self), &base_url_clone), labels: LabelsService::new(Weak::clone(weak_self), &base_url_clone), edges: EdgesService::new(Weak::clone(weak_self), &base_url_clone), + tenant: TenantService::new(Weak::clone(weak_self), &base_url_clone), http_client, } }); diff --git a/src/resources/mod.rs b/src/resources/mod.rs index f9df949..2eca052 100644 --- a/src/resources/mod.rs +++ b/src/resources/mod.rs @@ -213,6 +213,92 @@ impl ResourceService { self.execute_post_request::(url, form) .await } + + /// `GET /resources/export/{id}` — the whole connected graph component around one resource, as + /// a gzip-compressed file, in memory. + /// + /// The file names everything by external id, never by numeric id, so it imports into another + /// tenant or environment with [`import_graph`](Self::import_graph). A component over 2,000,000 + /// nodes or relationships is a 400 rather than a partial export. Use + /// [`export_graph_to_path`](Self::export_graph_to_path) for anything large. + pub async fn export_graph(&self, id: u64) -> Result, ResponseError> { + let path = format!("{}/export/{}", self.base_url, id); + let response = self.execute_get_stream_request(&path).await?; + let status = response.status(); + let bytes = response.bytes().await.map_err(|err| ResponseError { + status, + message: err.to_string(), + content_type: None, + })?; + Ok(bytes.to_vec()) + } + + /// [`export_graph`](Self::export_graph), streamed to `destination` without buffering the file + /// in memory. Returns the number of bytes written. The destination is created if missing and + /// truncated if it exists. + pub async fn export_graph_to_path( + &self, + id: u64, + destination: impl AsRef, + ) -> Result { + use tokio::io::AsyncWriteExt; + let path = format!("{}/export/{}", self.base_url, id); + let mut response = self.execute_get_stream_request(&path).await?; + let status = response.status(); + let io_error = |err: std::io::Error| ResponseError { + status, + message: err.to_string(), + content_type: None, + }; + let mut file = tokio::fs::File::create(destination.as_ref()) + .await + .map_err(io_error)?; + let mut written: u64 = 0; + while let Some(chunk) = response.chunk().await.map_err(|err| ResponseError { + status, + message: err.to_string(), + content_type: None, + })? { + file.write_all(&chunk).await.map_err(io_error)?; + written += chunk.len() as u64; + } + file.flush().await.map_err(io_error)?; + Ok(written) + } + + /// `POST /resources/import` — recreate the resources and relationships of a file produced by + /// [`export_graph`](Self::export_graph). + /// + /// What already exists is skipped rather than rejected — nodes by external id, relationships + /// by (from, to, type) — so re-importing into the source tenant is a no-op, and after a + /// failure the same file can simply be sent again: it fast-forwards through the segments + /// already committed. Timeseries are not created; missing ones are listed in + /// [`nodes_skipped_timeseries`](GraphImportResult::nodes_skipped_timeseries). Over 512 MB, or + /// 2,000,000 nodes or relationships, is a 413. + pub async fn import_graph( + &self, + file: impl Into, + ) -> Result { + let path = format!("{}/import", self.base_url); + self.execute_post_bytes_request(&path, file, "application/octet-stream") + .await + } + + /// [`import_graph`](Self::import_graph), streaming the file from disk. + pub async fn import_graph_from_path( + &self, + source: impl AsRef, + ) -> Result { + let file = tokio::fs::File::open(source.as_ref()).await.map_err(|e| { + ResponseError::bad_request(format!( + "failed to open '{}': {}", + source.as_ref().display(), + e + )) + })?; + let stream = tokio_util::codec::FramedRead::new(file, tokio_util::codec::BytesCodec::new()); + self.import_graph(reqwest::Body::wrap_stream(stream)).await + } } #[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] #[serde(rename_all = "camelCase")] @@ -641,3 +727,49 @@ impl FetchNearestResourcesForm { self } } + +/// Answer of [`ResourceService::import_graph`]. A bare object on the wire. +#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct GraphImportResult { + pub nodes_created: u64, + pub relations_created: u64, + /// Skipped because a node with the same external id already exists. + pub nodes_skipped_existing: u64, + /// Timeseries in the file that do not exist here. They cannot be created through the resource + /// api — create them through the timeseries api first, then import again. + #[serde(default)] + pub nodes_skipped_timeseries: Vec, + /// Already present, or an endpoint is unavailable. + pub relations_skipped: u64, + /// Nodes whose data set reference could not be resolved here and was dropped. + pub data_set_references_dropped: u64, + /// Transactions committed. The import streams in segments of 50,000 objects, each atomic. + pub segments: u64, + /// Naming-policy violations that were allowed through and recorded for review. + #[serde(default)] + pub warnings: Vec, +} + +impl crate::generic::DataWrapperDeserialization for GraphImportResult { + fn deserialize_and_set_status(body: &str, _status_code: u16) -> Result { + serde_json::from_str(body) + } +} + +/// A naming-policy violation that was allowed through and recorded for review. +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct PolicyWarning { + /// Position of the offending item in the submitted batch. + pub index: u32, + pub external_id: String, + /// External id of the policy that fired. + #[serde(default)] + pub policy: Option, + #[serde(default)] + pub message: Option, + /// A conforming alternative. Not applied — external ids are stored exactly as sent. + #[serde(default)] + pub suggestion: Option, +} diff --git a/src/resources/tests.rs b/src/resources/tests.rs index d157b83..45be52d 100644 --- a/src/resources/tests.rs +++ b/src/resources/tests.rs @@ -1115,3 +1115,131 @@ async fn update_echo_is_typed_per_node_type() -> Result<(), ResponseError> { fn_cleanup.disarm(); Ok(()) } + +/// Export a two-node component, then import the file back into the tenant it came from: every +/// node already exists, so the import is a no-op that skips rather than creates. +#[tokio::test] +async fn graph_export_then_import_into_the_source_tenant_is_a_no_op() -> Result<(), ResponseError> { + let api_service = create_api_service(); + let test_resources = create_test_resources(); + let root_ext = test_resources[0].external_id.clone(); + let child_ext = test_resources[1].external_id.clone(); + let relations = vec![RelForm::by_external_ids( + root_ext.clone(), + child_ext.clone(), + "flows_to", + )]; + let created = api_service + .resources + .create(test_resources, relations) + .await?; + let _cleanup = cleanup_resources(vec![child_ext.clone(), root_ext.clone()]); + let root_id = created + .nodes() + .unwrap() + .iter() + .find(|n| n.external_id() == root_ext) + .and_then(|n| n.id()) + .expect("create echoes the root's id"); + + // The export walks the graph projection, which lags the write. + let related = RelatedResourcesForm { + id: None, + external_id: Some(root_ext.clone()), + depth: -1, + relationship_types: None, + limit: 100, + excluded_labels: vec![], + }; + let network = poll_until( + || api_service.resources.fetch_related(&related), + |r| r.as_ref().map(|n| n.nodes().len() >= 2).unwrap_or(false), + ) + .await?; + assert!(network.nodes().len() >= 2, "graph projection did not catch up"); + + let file = api_service.resources.export_graph(root_id).await?; + assert_eq!(&file[..2], &[0x1f, 0x8b], "the export is gzip"); + + let result = api_service.resources.import_graph(file).await?; + assert_eq!(result.nodes_created, 0, "{result:?}"); + assert!(result.nodes_skipped_existing >= 2, "{result:?}"); + assert!(result.segments >= 1, "{result:?}"); + Ok(()) +} + +#[tokio::test] +async fn graph_export_of_an_unknown_id_is_a_404() { + let api_service = create_api_service(); + let err = api_service + .resources + .export_graph(u64::MAX / 2) + .await + .expect_err("no such resource"); + assert_eq!(err.status.as_u16(), 404); +} + +/// `root -> middle -> leaf`, with only the leaf carrying the end label: the nearest match is the +/// leaf, and the answer carries the path back to the start, so `middle` comes too. A relationship +/// type filter that matches no edge reaches nothing. +#[tokio::test] +async fn fetch_nearest_reaches_the_labelled_node_through_the_path() -> Result<(), ResponseError> { + use crate::tests::ids::TEST_LABEL; + + let api_service = create_api_service(); + let ids: Vec = ["nearest_root", "nearest_middle", "nearest_leaf"] + .iter() + .map(|kind| unique_id(kind)) + .collect(); + let node = |ext: &String, labels: Vec<&str>, is_root: bool| Resource { + is_root, + labels: Some(labels.into_iter().map(str::to_string).collect()), + ..{ + let mut r = Resource::new(); + r.external_id = ext.clone(); + r.name = format!("Rust SDK nearest probe {ext}"); + r + } + }; + let nodes = vec![ + node(&ids[0], vec!["ASSET"], true), + node(&ids[1], vec!["ASSET"], false), + node(&ids[2], vec!["ASSET", TEST_LABEL], false), + ]; + let relations = vec![ + RelForm::by_external_ids(ids[0].clone(), ids[1].clone(), "flows_to"), + RelForm::by_external_ids(ids[1].clone(), ids[2].clone(), "flows_to"), + ]; + let created = api_service.resources.create(nodes, relations).await?; + // Leaf first: the backend refuses to delete the start of an edge. + let _cleanup = cleanup_resources(ids.iter().rev().cloned().collect()); + let root_id = created + .nodes() + .unwrap() + .iter() + .find(|n| n.external_id() == ids[0]) + .and_then(|n| n.id()) + .expect("create echoes the root's id"); + + let form = FetchNearestResourcesForm { + end_labels: Some(vec![TEST_LABEL.to_string()]), + limit: Some(1), + ..FetchNearestResourcesForm::from_id(root_id) + }; + let has = |net: &ResourceNetwork, ext: &str| net.nodes().iter().any(|n| n.external_id() == ext); + let network = poll_until( + || api_service.resources.fetch_nearest(&form), + |r| r.as_ref().map(|n| has(n, &ids[2])).unwrap_or(false), + ) + .await?; + assert!(has(&network, &ids[2]), "the labelled leaf was not reached"); + assert!(has(&network, &ids[1]), "the path back to the start is part of the answer"); + + let unmatched = FetchNearestResourcesForm { + relationship_types: Some(vec![unique_id("no_such_type").to_uppercase()]), + ..form.clone() + }; + let none = api_service.resources.fetch_nearest(&unmatched).await?; + assert!(!has(&none, &ids[2]), "a relationship filter matching no edge still reached the leaf"); + Ok(()) +} diff --git a/src/tenant/mod.rs b/src/tenant/mod.rs new file mode 100644 index 0000000..e806f61 --- /dev/null +++ b/src/tenant/mod.rs @@ -0,0 +1,277 @@ +use crate::generic::{ApiServiceProvider, DataWrapperDeserialization}; +use crate::http::ResponseError; +use crate::ApiService; +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; +use std::sync::Weak; + +/// Client for the `/tenant` endpoints: what the tenant your token belongs to has switched on, and +/// the settings your organization administers for itself. +/// +/// Every answer here is a bare object, not an `items` envelope. +pub struct TenantService { + pub(crate) api_service: Weak, + base_url: String, +} + +impl ApiServiceProvider for TenantService { + fn api_service(&self) -> &Weak { + &self.api_service + } +} + +impl TenantService { + pub fn new(api_service: Weak, base_url: &String) -> Self { + TenantService { + api_service, + base_url: format!("{}/tenant", base_url), + } + } + + /// `GET /tenant/features` — which optional features are enabled for your tenant. A disabled + /// feature's endpoints may still exist and answer 404 or 403. + pub async fn features(&self) -> Result { + let path = &format!("{}/features", self.base_url); + self.execute_get_request(path, None::<&str>).await + } + + /// `GET /tenant/settings/permissions` — what the caller may read and write, per settings scope + /// (`"llm"`, …). Wildcard grants are already resolved, so every scope is listed by name. + /// + /// For gating a UI, not a security boundary: the settings endpoints enforce the same grants. + pub async fn settings_permissions( + &self, + ) -> Result, ResponseError> { + let path = &format!("{}/settings/permissions", self.base_url); + self.execute_get_request(path, None::<&str>).await + } + + /// `GET /tenant/settings/llm` — the model your organization's assistant runs on. The API key + /// is never returned; [`api_key_set`](TenantLlmSettings::api_key_set) says whether one is + /// stored. Needs the `llm` read grant (403 otherwise). + pub async fn llm_settings(&self) -> Result { + let path = &format!("{}/settings/llm", self.base_url); + self.execute_get_request(path, None::<&str>).await + } + + /// `PUT /tenant/settings/llm` — **replace** the model configuration, and answer it as stored. + /// + /// A replace, not a patch: a field left `None` is cleared. The one exception is + /// [`api_key`](TenantLlmSettingsForm::api_key), where `None` or empty keeps the stored + /// credential, so a form can save without retyping it. Needs the `llm` write grant. + pub async fn update_llm_settings( + &self, + form: &TenantLlmSettingsForm, + ) -> Result { + let path = &format!("{}/settings/llm", self.base_url); + self.execute_put_request(path, form).await + } +} + +/// Answer of [`TenantService::features`]. Each flag is `false` when the api leaves it unset. +#[derive(Debug, Serialize, Deserialize, Clone, Copy, Default, PartialEq, Eq)] +pub struct TenantFeatures { + #[serde(default)] + pub files: bool, + #[serde(default, deserialize_with = "null_as_false")] + pub policy: bool, + #[serde(default, deserialize_with = "null_as_false")] + pub streaming: bool, + #[serde(default, deserialize_with = "null_as_false")] + pub chat: bool, +} + +fn null_as_false<'de, D: serde::Deserializer<'de>>(deserializer: D) -> Result { + Ok(Option::::deserialize(deserializer)?.unwrap_or(false)) +} + +/// What the caller may do with one settings scope. +#[derive(Debug, Serialize, Deserialize, Clone, Copy, Default, PartialEq, Eq)] +pub struct SettingsPermission { + pub read: bool, + pub write: bool, +} + +/// Answer of [`TenantService::llm_settings`] and [`TenantService::update_llm_settings`]. +#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct TenantLlmSettings { + /// `anthropic` or `openai-compatible`. + #[serde(default)] + pub provider: Option, + #[serde(default)] + pub model: Option, + #[serde(default)] + pub base_url: Option, + #[serde(default)] + pub reasoning_effort: Option, + /// One of `low`, `medium`, `high`, `xhigh`, `max`. + #[serde(default)] + pub effort: Option, + #[serde(default)] + pub turn_timeout: Option, + #[serde(default)] + pub max_output_tokens: Option, + #[serde(default)] + pub max_iterations: Option, + #[serde(default)] + pub instructions: Option, + /// Whether a credential is stored. The credential itself is never returned. + #[serde(default)] + pub api_key_set: bool, + /// Whether this amounts to a model that can actually be called. `false` means your + /// organization has no assistant. + #[serde(default)] + pub configured: bool, +} + +/// Body of [`TenantService::update_llm_settings`]. +#[derive(Debug, Serialize, Deserialize, Clone, Default, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct TenantLlmSettingsForm { + #[serde(skip_serializing_if = "Option::is_none")] + pub provider: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub model: Option, + /// `None` or empty keeps the stored credential; only a non-blank value replaces it. + #[serde(skip_serializing_if = "Option::is_none")] + pub api_key: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub base_url: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub reasoning_effort: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub effort: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub turn_timeout: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub max_output_tokens: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub max_iterations: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub instructions: Option, +} + +impl From<&TenantLlmSettings> for TenantLlmSettingsForm { + /// The stored settings as a form that saves them unchanged — `api_key` stays `None`, which + /// keeps the stored credential. Edit the fields to change, then send. + fn from(value: &TenantLlmSettings) -> Self { + TenantLlmSettingsForm { + provider: value.provider.clone(), + model: value.model.clone(), + api_key: None, + base_url: value.base_url.clone(), + reasoning_effort: value.reasoning_effort.clone(), + effort: value.effort.clone(), + turn_timeout: value.turn_timeout.clone(), + max_output_tokens: value.max_output_tokens, + max_iterations: value.max_iterations, + instructions: value.instructions.clone(), + } + } +} + +impl DataWrapperDeserialization for TenantFeatures { + fn deserialize_and_set_status(body: &str, _status_code: u16) -> Result { + serde_json::from_str(body) + } +} + +impl DataWrapperDeserialization for TenantLlmSettings { + fn deserialize_and_set_status(body: &str, _status_code: u16) -> Result { + serde_json::from_str(body) + } +} + +impl DataWrapperDeserialization for HashMap { + fn deserialize_and_set_status(body: &str, _status_code: u16) -> Result { + serde_json::from_str(body) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::create_api_service; + + #[test] + fn features_read_null_flags_as_off() { + let features: TenantFeatures = + serde_json::from_str(r#"{"files":true,"policy":null,"streaming":true}"#).unwrap(); + assert!(features.files && features.streaming); + assert!(!features.policy && !features.chat); + } + + #[test] + fn the_llm_form_built_from_stored_settings_keeps_the_credential() { + let stored = TenantLlmSettings { + provider: Some("anthropic".into()), + api_key_set: true, + configured: true, + ..Default::default() + }; + let body = serde_json::to_value(TenantLlmSettingsForm::from(&stored)).unwrap(); + assert_eq!(body, serde_json::json!({"provider": "anthropic"})); + } + + #[tokio::test] + async fn features_and_settings_permissions_answer() { + let api = create_api_service(); + api.tenant.features().await.unwrap(); + let permissions = api.tenant.settings_permissions().await.unwrap(); + assert!(permissions.contains_key("llm"), "{permissions:?}"); + if permissions["llm"].read { + api.tenant.llm_settings().await.unwrap(); + } else { + let err = api.tenant.llm_settings().await.expect_err("no llm read grant"); + assert_eq!(err.status.as_u16(), 403); + } + } + + /// Without the `llm` write grant the `PUT` is a 403. With it, an empty form is a 400 naming + /// `provider` and `model` and writes nothing, and a configured model written back unchanged + /// answers exactly what was stored — the credential kept because `api_key` is left `None`. + /// The grant is checked before the form, so no branch can change the tenant's settings. + #[tokio::test] + async fn update_llm_settings_is_gated_validated_and_round_trips() { + let api = create_api_service(); + let permissions = api.tenant.settings_permissions().await.unwrap(); + let empty = TenantLlmSettingsForm::default(); + + if !permissions["llm"].write { + let err = api + .tenant + .update_llm_settings(&empty) + .await + .expect_err("no llm write grant"); + assert_eq!(err.status.as_u16(), 403); + println!("no llm write grant; the round trip is not exercised"); + return; + } + + let err = api + .tenant + .update_llm_settings(&empty) + .await + .expect_err("provider and model are required"); + assert_eq!(err.status.as_u16(), 400); + let fields: Vec = err + .problem() + .map(|p| p.fields().into_iter().filter_map(|f| f.field).collect()) + .unwrap_or_default(); + assert!(fields.contains(&"provider".to_string()), "{fields:?}"); + assert!(fields.contains(&"model".to_string()), "{fields:?}"); + + let stored = api.tenant.llm_settings().await.unwrap(); + if !stored.configured { + println!("no model configured; nothing to write back unchanged"); + return; + } + let saved = api + .tenant + .update_llm_settings(&TenantLlmSettingsForm::from(&stored)) + .await + .unwrap(); + assert_eq!(saved, stored); + } +} diff --git a/src/timeseries/datapoint_listen.rs b/src/timeseries/datapoint_listen.rs new file mode 100644 index 0000000..d32a587 --- /dev/null +++ b/src/timeseries/datapoint_listen.rs @@ -0,0 +1,455 @@ +//! Live datapoint tail over `ws(s):///timeseries/datapoints/listen`. +//! +//! Unlike [`SubscriptionListener`](crate::subscriptions::SubscriptionListener) there is no +//! subscription entity behind it and nothing to ack: each connection reads the firehose from +//! *latest*, non-durably, narrowed server-side to the tenant and the requested timeseries. Points +//! written while the socket is down are not replayed. + +use futures::{SinkExt, StreamExt}; +use serde::{Deserialize, Serialize}; +use std::collections::VecDeque; +use std::sync::Weak; +use std::time::Duration; +use tokio::net::TcpStream; +use tokio_tungstenite::tungstenite::client::IntoClientRequest; +use tokio_tungstenite::tungstenite::http; +use tokio_tungstenite::tungstenite::protocol::frame::coding::CloseCode; +use tokio_tungstenite::tungstenite::Message; +use tokio_tungstenite::{connect_async, MaybeTlsStream, WebSocketStream}; + +use crate::subscriptions::ListenError; +use crate::ApiService; + +/// Offered alongside the bearer element so the server has a subprotocol to echo that is not the +/// credential itself. +const NEGOTIATED_SUBPROTOCOL: &str = "datahub.v1"; +const BEARER_SUBPROTOCOL_PREFIX: &str = "datahub.bearer."; + +const RECONNECT_INITIAL_BACKOFF: Duration = Duration::from_millis(500); +const RECONNECT_MAX_BACKOFF: Duration = Duration::from_secs(30); +const RECONNECT_MAX_RETRIES: u32 = 8; + +/// One point delivered by [`DatapointListener::next`]. `value` is a string for every value type, +/// as on the other ingest and listen paths. +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct LiveDatapoint { + pub external_id: String, + /// Upper case, e.g. `FLOAT`. + #[serde(default)] + pub value_type: Option, + pub timestamp: String, + pub value: String, +} + +#[derive(Debug, Deserialize)] +#[serde(untagged)] +enum ServerFrame { + Datapoints { + datapoints: Vec, + }, + Error { + #[allow(dead_code)] + error: bool, + #[serde(default)] + reason: Option, + #[serde(default)] + scope: Option, + #[serde(default)] + limit: Option, + #[serde(default)] + message: Option, + }, +} + +pub(crate) fn decode_frame(text: &str) -> Result, ListenError> { + match serde_json::from_str(text).map_err(|e| ListenError::Deserialize(e.to_string()))? { + ServerFrame::Datapoints { datapoints } => Ok(datapoints), + ServerFrame::Error { + scope: Some(scope), + limit: Some(limit), + message, + reason, + .. + } => Err(ListenError::ConnectionLimit { + scope, + limit, + message: message + .or(reason) + .unwrap_or_else(|| "connection refused".to_string()), + }), + ServerFrame::Error { + message, reason, .. + } => Err(ListenError::WebSocket( + message + .or(reason) + .unwrap_or_else(|| "server reported an error".to_string()), + )), + } +} + +/// Live tail of the datapoints written to a set of timeseries, by external id. +/// +/// Drive it by calling [`next`](Self::next) in a loop; change what is streamed with +/// [`subscribe`](Self::subscribe) / [`unsubscribe`](Self::unsubscribe) / +/// [`set_timeseries`](Self::set_timeseries). The server narrows every change to the timeseries +/// whose data set the caller may read, **silently** — an unreadable or unknown external id is +/// dropped, not reported. +/// +/// A dropped connection is re-established by `next` with a fresh token and the current interest +/// set; points written in the gap are lost. A refusal — a bad token, a missing role, the connection +/// limit — is returned rather than retried. +pub struct DatapointListener { + ws: WebSocketStream>, + buffered: VecDeque, + api_service: Weak, + ws_url: String, + interest: Vec, +} + +impl DatapointListener { + pub(crate) async fn connect( + api_service: Weak, + timeseries_base_url: &str, + interest: Vec, + ) -> Result { + let ws_url = build_ws_url(timeseries_base_url)?; + let ws = Self::open(&api_service, &ws_url, &interest).await?; + Ok(DatapointListener { + ws, + buffered: VecDeque::new(), + api_service, + ws_url, + interest, + }) + } + + async fn open( + api_service: &Weak, + ws_url: &str, + interest: &[String], + ) -> Result>, ListenError> { + let service = api_service + .upgrade() + .ok_or_else(|| ListenError::Request("api service has been dropped".to_string()))?; + let token = service + .config + .get_api_token() + .await + .map_err(|e| ListenError::Request(format!("failed to get api token: {}", e)))?; + + let mut url = + reqwest::Url::parse(ws_url).map_err(|e| ListenError::Request(e.to_string()))?; + if !interest.is_empty() { + url.query_pairs_mut() + .append_pair("externalIds", &interest.join(",")); + } + let mut request = url + .as_str() + .into_client_request() + .map_err(|e| ListenError::Request(e.to_string()))?; + // A browser cannot set Authorization on a handshake, so the api reads the token from the + // offered subprotocols instead and ignores the header. No space after the comma: + // tungstenite splits the offer on "," without trimming, and would then reject the echoed + // `datahub.v1` as one it never offered. + let protocols: http::HeaderValue = + format!("{BEARER_SUBPROTOCOL_PREFIX}{token},{NEGOTIATED_SUBPROTOCOL}") + .parse() + .map_err(|e: http::header::InvalidHeaderValue| { + ListenError::Request(e.to_string()) + })?; + request + .headers_mut() + .insert(http::header::SEC_WEBSOCKET_PROTOCOL, protocols); + + let (ws, _response) = connect_async(request) + .await + .map_err(|e| ListenError::Handshake(e.to_string()))?; + Ok(ws) + } + + async fn reconnect(&mut self) -> Result<(), ListenError> { + let mut delay = RECONNECT_INITIAL_BACKOFF; + let mut last_err = ListenError::WebSocket("connection lost".to_string()); + for _ in 0..RECONNECT_MAX_RETRIES { + tokio::time::sleep(delay).await; + match Self::open(&self.api_service, &self.ws_url, &self.interest).await { + Ok(ws) => { + self.ws = ws; + return Ok(()); + } + Err(e) => { + last_err = e; + delay = (delay * 2).min(RECONNECT_MAX_BACKOFF); + } + } + } + Err(last_err) + } + + /// Wait for the next datapoint. Reconnects transparently when the connection drops. Returns + /// `Some(Err(_))` when a reconnect ultimately fails, a frame cannot be decoded, or the server + /// refuses the connection. + pub async fn next(&mut self) -> Option> { + loop { + if let Some(point) = self.buffered.pop_front() { + return Some(Ok(point)); + } + let frame = match self.ws.next().await { + Some(Ok(f)) => f, + None | Some(Err(_)) => match self.reconnect().await { + Ok(()) => continue, + Err(e) => return Some(Err(e)), + }, + }; + match frame { + Message::Text(text) => match decode_frame(&text) { + Ok(points) => self.buffered.extend(points), + Err(e) => return Some(Err(e)), + }, + // The api authenticates after the 101, so a bad token, a missing role or a missing + // tenant arrives as a policy-violation close. Reconnecting would repeat it. + Message::Close(Some(close)) if close.code == CloseCode::Policy => { + return Some(Err(ListenError::Handshake(close.reason.to_string()))); + } + Message::Close(_) => match self.reconnect().await { + Ok(()) => continue, + Err(e) => return Some(Err(e)), + }, + _ => continue, + } + } + } + + /// Add timeseries to the live set. + pub async fn subscribe>(&mut self, external_ids: &[S]) -> Result<(), ListenError> { + for id in external_ids { + let id = id.as_ref().to_string(); + if !self.interest.contains(&id) { + self.interest.push(id); + } + } + self.send_interest("subscribe", external_ids).await + } + + /// Remove timeseries from the live set. + pub async fn unsubscribe>( + &mut self, + external_ids: &[S], + ) -> Result<(), ListenError> { + let removing: Vec = external_ids.iter().map(|s| s.as_ref().to_string()).collect(); + self.interest.retain(|id| !removing.contains(id)); + self.send_interest("unsubscribe", external_ids).await + } + + /// Replace the whole live set. + pub async fn set_timeseries>( + &mut self, + external_ids: &[S], + ) -> Result<(), ListenError> { + self.interest = external_ids.iter().map(|s| s.as_ref().to_string()).collect(); + self.send_interest("set", external_ids).await + } + + async fn send_interest>( + &mut self, + action: &str, + external_ids: &[S], + ) -> Result<(), ListenError> { + let ids: Vec<&str> = external_ids.iter().map(|s| s.as_ref()).collect(); + let frame = serde_json::to_string(&serde_json::json!({ + "action": action, + "externalIds": ids, + }))?; + self.ws + .send(Message::Text(frame.into())) + .await + .map_err(|e| ListenError::WebSocket(e.to_string())) + } + + /// Send a Close frame and drain remaining frames until the peer closes its side. + pub async fn close(mut self) -> Result<(), ListenError> { + let _ = self.ws.close(None).await; + while let Some(frame) = self.ws.next().await { + if frame.is_err() { + break; + } + } + Ok(()) + } +} + +/// `http(s):///timeseries` to `ws(s):///timeseries/datapoints/listen`. +pub(crate) fn build_ws_url(timeseries_base_url: &str) -> Result { + let ws_base = if let Some(rest) = timeseries_base_url.strip_prefix("https://") { + format!("wss://{}", rest) + } else if let Some(rest) = timeseries_base_url.strip_prefix("http://") { + format!("ws://{}", rest) + } else { + return Err(ListenError::Request(format!( + "base_url must start with http:// or https://, got {}", + timeseries_base_url + ))); + }; + Ok(format!( + "{}/datapoints/listen", + ws_base.trim_end_matches('/') + )) +} + +#[cfg(test)] +mod tests { + use super::*; + + use crate::datahub::DataHubConfig; + use crate::ApiService; + use tokio::net::TcpListener; + use tokio_tungstenite::tungstenite::handshake::server::{Request, Response}; + use tokio_tungstenite::tungstenite::protocol::CloseFrame; + + /// A server that authenticates the way the api's handler does — the token from the offered + /// subprotocols, `datahub.v1` echoed — then runs `session` on the socket. + async fn fake_listen_endpoint(session: F) -> (std::sync::Arc, tokio::task::JoinHandle>) + where + F: FnOnce(WebSocketStream) -> Fut + Send + 'static, + Fut: std::future::Future + Send, + { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let handle = tokio::spawn(async move { + let (socket, _) = listener.accept().await.unwrap(); + let mut offered = None; + let mut query = None; + let ws = tokio_tungstenite::accept_hdr_async(socket, |req: &Request, mut resp: Response| { + assert_eq!(req.uri().path(), "/timeseries/datapoints/listen"); + query = req.uri().query().map(str::to_string); + offered = req + .headers() + .get(http::header::SEC_WEBSOCKET_PROTOCOL) + .map(|v| v.to_str().unwrap().to_string()); + resp.headers_mut().insert( + http::header::SEC_WEBSOCKET_PROTOCOL, + http::HeaderValue::from_static(NEGOTIATED_SUBPROTOCOL), + ); + Ok(resp) + }) + .await + .unwrap(); + let offered = offered.unwrap_or_default(); + let bearer = offered + .split(',') + .map(str::trim) + .find_map(|p| p.strip_prefix(BEARER_SUBPROTOCOL_PREFIX)) + .map(str::to_string); + assert_eq!(bearer.as_deref(), Some("static-token"), "offered: {offered}"); + session(ws).await; + query + }); + let config = DataHubConfig::from_vars( + format!("http://{addr}"), + Some("static-token".to_string()), + None, + None, + None, + None, + ); + (ApiService::new(config), handle) + } + + #[tokio::test] + async fn the_token_rides_in_the_subprotocol_and_points_arrive() { + let (api, server) = fake_listen_endpoint(|mut ws| async move { + ws.send(Message::Text( + r#"{"datapoints":[{"externalId":"a","valueType":"FLOAT","timestamp":"2026-01-01T00:00:00Z","value":"1.5"}]}"#.into(), + )) + .await + .unwrap(); + let Some(Ok(Message::Text(change))) = ws.next().await else { + panic!("expected the interest change"); + }; + let change: serde_json::Value = serde_json::from_str(&change).unwrap(); + assert_eq!(change, serde_json::json!({"action": "subscribe", "externalIds": ["b"]})); + let _ = ws.close(None).await; + }) + .await; + + let mut listener = api.time_series.listen_datapoints(&["a", "x y"]).await.unwrap(); + let point = listener.next().await.unwrap().unwrap(); + assert_eq!(point.external_id, "a"); + listener.subscribe(&["b"]).await.unwrap(); + let query = server.await.unwrap(); + assert_eq!(query.as_deref(), Some("externalIds=a%2Cx+y")); + } + + /// The api authenticates after the 101 and refuses with a policy-violation close; that is + /// returned, not retried. + #[tokio::test] + async fn a_policy_close_is_returned_rather_than_retried() { + let (api, server) = fake_listen_endpoint(|mut ws| async move { + let _ = ws + .close(Some(CloseFrame { + code: CloseCode::Policy, + reason: "Invalid access token".into(), + })) + .await; + }) + .await; + + let mut listener = api.time_series.listen_datapoints::<&str>(&[]).await.unwrap(); + let outcome = tokio::time::timeout(Duration::from_secs(2), listener.next()) + .await + .expect("no reconnect loop"); + match outcome { + Some(Err(ListenError::Handshake(reason))) => assert_eq!(reason, "Invalid access token"), + other => panic!("expected the refusal, got {other:?}"), + } + assert_eq!(server.await.unwrap(), None, "an empty interest set sends no query"); + } + + #[test] + fn ws_url_swaps_the_scheme_and_appends_the_listen_path() { + assert_eq!( + build_ws_url("https://api.example.com/timeseries").unwrap(), + "wss://api.example.com/timeseries/datapoints/listen" + ); + assert_eq!( + build_ws_url("http://localhost:8081/timeseries").unwrap(), + "ws://localhost:8081/timeseries/datapoints/listen" + ); + assert!(build_ws_url("ftp://x/timeseries").is_err()); + } + + #[test] + fn a_datapoints_frame_decodes_every_point() { + let points = decode_frame( + r#"{"datapoints":[ + {"externalId":"a","valueType":"FLOAT","timestamp":"2026-01-01T00:00:00Z","value":"1.5"}, + {"externalId":"b","valueType":"TEXT","timestamp":"2026-01-01T00:00:01Z","value":"on"} + ]}"#, + ) + .unwrap(); + assert_eq!(points.len(), 2); + assert_eq!(points[0].external_id, "a"); + assert_eq!(points[1].value, "on"); + } + + #[test] + fn the_limit_refusal_names_its_scope_and_cap() { + let err = decode_frame( + r#"{"error":true,"reason":"websocket-limit-reached","scope":"user","limit":5,"message":"close one"}"#, + ) + .unwrap_err(); + match err { + ListenError::ConnectionLimit { + scope, + limit, + message, + } => { + assert_eq!(scope, "user"); + assert_eq!(limit, 5); + assert_eq!(message, "close one"); + } + other => panic!("expected ConnectionLimit, got {other:?}"), + } + } +} diff --git a/src/timeseries/mod.rs b/src/timeseries/mod.rs index ee9752d..dba329c 100644 --- a/src/timeseries/mod.rs +++ b/src/timeseries/mod.rs @@ -1,7 +1,9 @@ pub mod binary; +pub mod datapoint_listen; mod test; pub use binary::{BinaryIngestOptions, DatapointValueType, Frame, FrameWriter, ResolvedSeries}; +pub use datapoint_listen::{DatapointListener, LiveDatapoint}; use crate::buffer::DurableSpool; use crate::datahub::DataHubConfig; @@ -71,6 +73,58 @@ impl TimeSeriesService { .await } + /// `GET /timeseries/{id}` — one series by its numeric id. + /// + /// **404 does not mean the id is free.** A series the caller may not read is reported as + /// missing rather than forbidden. Unlike [`by_ids`](Self::by_ids), which omits what it cannot + /// find, a miss here is an error carrying the `not-found` problem type. + pub async fn get_by_id(&self, id: u64) -> Result, ResponseError> { + let path = &format!("{}/{}", self.base_url, id); + self.execute_get_request::, ()>(path, None) + .await + } + + /// Open a live tail of the datapoints written to these timeseries, by external id — the + /// `/timeseries/datapoints/listen` WebSocket. May be empty; add series later with + /// [`DatapointListener::subscribe`]. + /// + /// No subscription entity is involved and nothing is durable: the stream starts at *latest*, + /// and ids the caller cannot read are dropped silently. For at-least-once delivery use + /// [`SubscriptionsService::listen`](crate::subscriptions::SubscriptionsService::listen). + pub async fn listen_datapoints>( + &self, + external_ids: &[S], + ) -> Result { + let interest = external_ids.iter().map(|s| s.as_ref().to_string()).collect(); + DatapointListener::connect(self.api_service.clone(), &self.base_url, interest).await + } + + /// `GET /timeseries/recommend-value-type/{unitExternalId}` — the value type that compresses + /// best in ClickHouse for a unit while still representing it faithfully. + /// + /// Advice, not a constraint, and a hard-coded heuristic. An unknown unit is **not** an error: + /// it answers the generic compact default with + /// [`recognized`](ValueTypeRecommendation::recognized) `false`. + pub async fn recommend_value_type( + &self, + unit_external_id: &str, + ) -> Result { + let mut url = reqwest::Url::parse(&self.base_url).map_err(|e| ResponseError { + status: reqwest::StatusCode::BAD_REQUEST, + message: e.to_string(), + content_type: None, + })?; + url.path_segments_mut() + .map_err(|_| ResponseError { + status: reqwest::StatusCode::BAD_REQUEST, + message: format!("base url {} cannot carry a path", self.base_url), + content_type: None, + })? + .push("recommend-value-type") + .push(unit_external_id); + self.execute_get_request(url.as_str(), None::<&str>).await + } + pub async fn create( &self, json: &DataWrapper, @@ -952,3 +1006,23 @@ mod unbuffered_insert_tests { assert!(observed.requests.load(Ordering::SeqCst) < 5); } } + +/// Answer of [`TimeSeriesService::recommend_value_type`]. A bare object on the wire, not an +/// `items` envelope. +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq)] +#[serde(rename_all = "camelCase")] +pub struct ValueTypeRecommendation { + /// The unit external id the recommendation was made for, echoed from the request. + pub unit_external_id: String, + /// One of `BIGINT`, `FLOAT`, `FLOAT32`, `NUMERIC`, `DECIMAL32`, `TEXT` or `MIXED`. + pub recommended_value_type: String, + pub reason: String, + /// `false` when the unit matched nothing specific and the generic default came back. + pub recognized: bool, +} + +impl crate::generic::DataWrapperDeserialization for ValueTypeRecommendation { + fn deserialize_and_set_status(body: &str, _status_code: u16) -> Result { + serde_json::from_str(body) + } +} diff --git a/src/timeseries/test.rs b/src/timeseries/test.rs index d895c51..3332f93 100644 --- a/src/timeseries/test.rs +++ b/src/timeseries/test.rs @@ -1806,3 +1806,101 @@ fn timeseries_update_matches_the_server_field_set() { }; assert_eq!(cleared["source"], serde_json::json!({"set": null, "setNull": true})); } + +#[cfg(test)] +mod single_reads_and_live_tail { + use crate::create_api_service; + use crate::generic::{DataWrapper, DatapointString, DatapointsCollection}; + use crate::tests::cleanup::cleanup_timeseries; + use crate::tests::ids::unique_id; + use crate::timeseries::TimeSeries; + use chrono::Utc; + use std::time::Duration; + + async fn create_series(ext_id: &str) -> u64 { + let api = create_api_service(); + let mut ts = TimeSeries::new(ext_id, "Rust SDK single-read probe"); + ts.unit = Some("celsius".to_string()); + let created = api + .time_series + .create(&DataWrapper::from_vec(vec![ts])) + .await + .expect("create the probe series"); + created.get_items()[0].id.expect("create echoes the id") + } + + #[tokio::test] + async fn get_by_id_reads_one_series_and_404s_a_miss() { + let api = create_api_service(); + let ext_id = unique_id("ts"); + let id = create_series(&ext_id).await; + let _cleanup = cleanup_timeseries(vec![ext_id.clone()]); + + let got = api.time_series.get_by_id(id).await.unwrap(); + assert_eq!(got.get_items().len(), 1); + assert_eq!(got.get_items()[0].external_id, ext_id); + + let err = api + .time_series + .get_by_id(u64::MAX / 2) + .await + .expect_err("an unknown id is a 404"); + assert_eq!(err.status.as_u16(), 404); + assert_eq!(err.problem_slug().as_deref(), Some("not-found")); + } + + #[tokio::test] + async fn recommend_value_type_answers_known_and_unknown_units() { + let api = create_api_service(); + let known = api + .time_series + .recommend_value_type("temperature_deg_c") + .await + .unwrap(); + assert_eq!(known.unit_external_id, "temperature_deg_c"); + assert!(!known.recommended_value_type.is_empty()); + + let unknown_unit = unique_id("unit"); + let unknown = api + .time_series + .recommend_value_type(&unknown_unit) + .await + .unwrap(); + assert!(!unknown.recognized, "an unknown unit is the generic default, not an error"); + } + + /// Points written after the listener connects arrive on it; nothing is replayed from before. + #[tokio::test] + async fn listen_datapoints_delivers_points_written_after_connecting() { + let api = create_api_service(); + let ext_id = unique_id("ts"); + create_series(&ext_id).await; + let _cleanup = cleanup_timeseries(vec![ext_id.clone()]); + + let mut listener = api + .time_series + .listen_datapoints(&[ext_id.as_str()]) + .await + .expect("open the datapoint tail"); + // The consumer reads from latest; give it a moment to attach before writing. + tokio::time::sleep(Duration::from_secs(2)).await; + + let mut collection = DatapointsCollection::from_external_id(&ext_id); + collection + .datapoints + .push(DatapointString::from_datetime(Utc::now(), "21.5")); + api.time_series + .insert_datapoints(&mut DataWrapper::from_vec(vec![collection])) + .await + .expect("write a point"); + + let point = tokio::time::timeout(Duration::from_secs(30), listener.next()) + .await + .expect("a point within 30s") + .expect("the stream is open") + .expect("a decodable point"); + assert_eq!(point.external_id, ext_id); + assert_eq!(point.value.parse::().unwrap(), 21.5); + listener.close().await.unwrap(); + } +}