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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
42 changes: 33 additions & 9 deletions AGENTS.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand All @@ -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<u64>`, 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`)
Expand Down Expand Up @@ -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 `<epoch_millis>\t<json>`; the spool is content-agnostic.

### Live datapoint tail (`src/timeseries/datapoint_listen.rs`)

`TimeSeriesService::listen_datapoints` opens `ws(s)://<host>/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.<jwt>,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<DatapointString>` 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.
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading