diff --git a/AGENTS.md b/AGENTS.md index da76ca8..f6c68ab 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -86,7 +86,7 @@ This crate is a thin async HTTP SDK around a DataHub-style REST API. Entry point `update` echoes a typed `Node` rather than a flat `Resource`. - `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` is a strict subset of what `/filter` now accepts: it carries no `cursor`, and `SubscriptionFilter` has only `timeseries`, not the `id`, `externalId`, `name`, `createdTime` and `lastUpdatedTime` criteria the api's filter grew beside it. +- `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 diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index ccde6d0..c88c205 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -100,8 +100,12 @@ class DataHubException(Exception): All three are `None` when the API answered with something that is not a problem document — an empty 401, a stack trace, or plain text. + + A subscription listener raises it too, for a failure on its WebSocket. There + is no HTTP response behind that, so `status_code` is `None` along with the + problem attributes. Using a listener after `close()` raises `ValueError`. """ - status_code: int + status_code: int | None message: str problem: dict | None problem_type: str | None @@ -1947,6 +1951,9 @@ class INode: class FileUpload: + """A local file to upload. Constructing one reads the file's metadata, so a + missing path raises `FileNotFoundError` and a directory `IsADirectoryError`.""" + def __init__( self, path: str, @@ -2075,36 +2082,29 @@ SubscriptionTimeseriesId = Union[TimeSeries, IdCollection, int, str] class SubscriptionFilter: - def __init__(self, timeseries: list[SubscriptionTimeseriesId] | None = None) -> None: ... - @property - def timeseries(self) -> list[IdCollection]: ... - + """AND-combined criteria for ``subscriptions.filter``. -class DataSort: + ``external_id`` and ``name`` are pattern lists — see ``PatternList``. ``timeseries`` matches + subscriptions bound to at least one of the given timeseries. A subscription is not a node, so + there is no ``source``, ``labels`` or ``metadata``. + """ def __init__( self, - property: list[str] | None = None, - order: str | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + timeseries: list[SubscriptionTimeseriesId] | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, ) -> None: ... @property - def property(self) -> list[str]: ... + def id(self) -> list[int] | None: ... @property - def order(self) -> str | None: ... - - -class SubscriptionFilterForm: - def __init__( - self, - filter: SubscriptionFilter | None = None, - limit: int | None = None, - sort: DataSort | None = None, - ) -> None: ... - @property - def filter(self) -> SubscriptionFilter: ... + def external_id(self) -> list[str] | None: ... @property - def limit(self) -> int: ... + def name(self) -> list[str] | None: ... @property - def sort(self) -> DataSort | None: ... + def timeseries(self) -> list[IdCollection]: ... class EventAction: @@ -2196,11 +2196,23 @@ class SubscriptionsServiceSync: def list(self, limit: int | None = None) -> list[Subscription]: ... def filter( self, - form: SubscriptionFilterForm | None = None, + *, + filter: SubscriptionFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, timeseries: list[SubscriptionTimeseriesId] | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, limit: int | None = None, - sort: DataSort | None = None, - ) -> list[Subscription]: ... + 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``. Sortable by ``id``, ``externalId``, ``name``, ``createdTime`` and + ``lastUpdatedTime``; the default is newest created first. + """ def delete(self, input: list[SubscriptionIdentifiable]) -> None: ... def listen(self, subscription_external_ids: list[str]) -> SubscriptionListener: ... @@ -2210,11 +2222,23 @@ class SubscriptionsServiceAsync: async def list(self, limit: int | None = None) -> list[Subscription]: ... async def filter( self, - form: SubscriptionFilterForm | None = None, + *, + filter: SubscriptionFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, timeseries: list[SubscriptionTimeseriesId] | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, limit: int | None = None, - sort: DataSort | None = None, - ) -> list[Subscription]: ... + 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``. Sortable by ``id``, ``externalId``, ``name``, ``createdTime`` and + ``lastUpdatedTime``; the default is newest created first. + """ async def delete(self, input: list[SubscriptionIdentifiable]) -> None: ... async def listen(self, subscription_external_ids: list[str]) -> SubscriptionListenerAsync: ... diff --git a/datahub_python_bindings/src/datasets/mod.rs b/datahub_python_bindings/src/datasets/mod.rs index edb4c7b..22ea024 100644 --- a/datahub_python_bindings/src/datasets/mod.rs +++ b/datahub_python_bindings/src/datasets/mod.rs @@ -484,7 +484,7 @@ impl PyDatasetFilter { /// /// Shared by the sync and async services so the accepted keywords cannot drift apart between them. /// -/// `limit` defaults to the server's 100 and may not exceed 10000 — above that the request is +/// `limit` defaults to 1000 and may not exceed 10000 — above that the request is /// rejected. #[allow(clippy::too_many_arguments)] pub fn dataset_filter_form( diff --git a/datahub_python_bindings/src/events/mod.rs b/datahub_python_bindings/src/events/mod.rs index 7911c1b..190ed05 100644 --- a/datahub_python_bindings/src/events/mod.rs +++ b/datahub_python_bindings/src/events/mod.rs @@ -120,7 +120,9 @@ pub fn event_filter_form( from_keywords, any_keyword, )?); - form.set_limit(limit.unwrap_or(100)); + if let Some(limit) = limit { + form.set_limit(limit); + } // A bare string is a one-element list here as it is on every other filter field; only one // property is used either way. if let Some(property) = sort_by { diff --git a/datahub_python_bindings/src/files/mod.rs b/datahub_python_bindings/src/files/mod.rs index df7f9f4..f4cfc39 100644 --- a/datahub_python_bindings/src/files/mod.rs +++ b/datahub_python_bindings/src/files/mod.rs @@ -304,7 +304,7 @@ impl PyFileUpload { data_set_id: Option, related_resources: Option>, ) -> PyResult { - let mut file_upload = FileUpload::new(path); + let mut file_upload = FileUpload::new(path)?; if let Some(external_id) = external_id { file_upload.external_id = external_id.to_string(); } @@ -334,7 +334,7 @@ impl PyFileUpload { #[classmethod] pub fn from_path(_py: Py, path: &str) -> PyResult { Ok(Self { - inner: FileUpload::new(path), + inner: FileUpload::new(path)?, }) } #[classmethod] @@ -344,7 +344,7 @@ impl PyFileUpload { destination_path: &str, ) -> PyResult { Ok(Self { - inner: FileUpload::new_with_destination_path(path, destination_path), + inner: FileUpload::new_with_destination_path(path, destination_path)?, }) } #[getter] diff --git a/datahub_python_bindings/src/lib.rs b/datahub_python_bindings/src/lib.rs index 5732966..a9a81ba 100644 --- a/datahub_python_bindings/src/lib.rs +++ b/datahub_python_bindings/src/lib.rs @@ -68,6 +68,22 @@ create_exception!( /// Convert an SDK `ResponseError` into a `DataHubException` that exposes the HTTP /// `status_code` and `message` as attributes, so Python callers can branch on the /// status code (e.g. `except DataHubException as e: if e.status_code == 409: ...`). +/// Convert a WebSocket listener failure into a `DataHubException`. There is no HTTP response behind +/// one, so `status_code` and the problem attributes are `None`, set anyway so an `except` block can +/// read them without `hasattr`. +pub(crate) fn listen_err(e: intellistream_datahub_sdk::ListenError) -> PyErr { + Python::attach(|py| { + let err = DataHubException::new_err(e.to_string()); + let value = err.value(py); + let _ = value.setattr("status_code", py.None()); + let _ = value.setattr("message", e.to_string()); + let _ = value.setattr("problem", py.None()); + let _ = value.setattr("problem_type", py.None()); + let _ = value.setattr("problem_slug", py.None()); + err + }) +} + /// Map a "not found" into Python's `None`. /// /// Every single-resource `GET` in the API answers an unknown id with **404** (batch `/byids` diff --git a/datahub_python_bindings/src/subscriptions/async_service.rs b/datahub_python_bindings/src/subscriptions/async_service.rs index 0325156..c4b7570 100644 --- a/datahub_python_bindings/src/subscriptions/async_service.rs +++ b/datahub_python_bindings/src/subscriptions/async_service.rs @@ -1,13 +1,11 @@ use crate::subscriptions::listener::{PySubscriptionListenerAsync, shared_listener}; -use crate::subscriptions::sync_service::build_filter_form; use crate::subscriptions::{ - PyDataSort, PySubscription, PySubscriptionFilterForm, SubscriptionIdentifyable, - SubscriptionTimeseriesId, + PySubscription, PySubscriptionFilter, SubscriptionIdentifyable, SubscriptionTimeseriesId, + subscription_filter_form, }; use intellistream_datahub_sdk::ApiService; use intellistream_datahub_sdk::generic::IdAndExtId; use intellistream_datahub_sdk::subscriptions::Subscription; -use pyo3::exceptions::PyException; use pyo3::prelude::*; use pyo3_async_runtimes::tokio::future_into_py; use std::sync::Arc; @@ -62,30 +60,41 @@ impl PySubscriptionsServiceAsync { }) } - /// Subscriptions matching every criterion on the filter. - #[pyo3(signature=(form=None, *, timeseries=None, limit=None, sort=None))] + /// Subscriptions matching every criterion on the filter, newest first. + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, timeseries=None, + created_time=None, last_updated_time=None, limit=None, sort_by=None, + sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] fn filter<'py>( &self, py: Python<'py>, - form: Option, + filter: Option, + id: Option>, + external_id: Option, + name: Option, timeseries: Option>, - limit: Option, - sort: Option, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, ) -> PyResult> { - let form = build_filter_form(form, timeseries, limit, sort)?; + let form = subscription_filter_form( + filter, id, external_id, name, timeseries, created_time, last_updated_time, limit, + sort_by, sort_order, cursor, + )?; let service = self.api_service.clone(); future_into_py(py, async move { let result = service .subscriptions .filter(&form) .await - .map_err(|e| crate::datahub_err(e))?; - Ok(result - .get_items() - .iter() - .cloned() - .map(PySubscription::from) - .collect::>()) + .map_err(crate::datahub_err)?; + let next_cursor = result.next_cursor().map(str::to_string); + let items: Vec = + result.get_items().iter().cloned().map(PySubscription::from).collect(); + Python::attach(|py| crate::PyPage::new(py, items, next_cursor)) }) } @@ -117,7 +126,7 @@ impl PySubscriptionsServiceAsync { .subscriptions .listen(&subscription_external_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(PySubscriptionListenerAsync { listener: shared_listener(listener), }) diff --git a/datahub_python_bindings/src/subscriptions/listener.rs b/datahub_python_bindings/src/subscriptions/listener.rs index 2090cfa..4ed78e2 100644 --- a/datahub_python_bindings/src/subscriptions/listener.rs +++ b/datahub_python_bindings/src/subscriptions/listener.rs @@ -1,6 +1,6 @@ use crate::subscriptions::PySubscriptionMessage; use intellistream_datahub_sdk::subscriptions::SubscriptionListener; -use pyo3::exceptions::{PyException, PyStopAsyncIteration, PyStopIteration}; +use pyo3::exceptions::{PyStopAsyncIteration, PyStopIteration, PyValueError}; use pyo3::prelude::*; use pyo3_async_runtimes::tokio::future_into_py; use std::sync::Arc; @@ -31,10 +31,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; match l.next().await { Some(Ok(msg)) => Ok(PySubscriptionMessage::from(msg)), - Some(Err(e)) => Err(PyException::new_err(e.to_string())), + Some(Err(e)) => Err(crate::listen_err(e)), None => Err(PyStopIteration::new_err(())), } }) @@ -52,10 +52,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; match l.next().await { Some(Ok(msg)) => Ok(Some(PySubscriptionMessage::from(msg))), - Some(Err(e)) => Err(PyException::new_err(e.to_string())), + Some(Err(e)) => Err(crate::listen_err(e)), None => Ok(None), } }) @@ -70,10 +70,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.ack(&message_ids) .await - .map_err(|e| PyException::new_err(e.to_string())) + .map_err(crate::listen_err) }) }) } @@ -86,10 +86,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.nack(&message_ids) .await - .map_err(|e| PyException::new_err(e.to_string())) + .map_err(crate::listen_err) }) }) } @@ -102,10 +102,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.subscribe(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string())) + .map_err(crate::listen_err) }) }) } @@ -118,10 +118,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.unsubscribe(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string())) + .map_err(crate::listen_err) }) }) } @@ -134,10 +134,10 @@ impl PySubscriptionListener { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.set_subscriptions(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string())) + .map_err(crate::listen_err) }) }) } @@ -151,7 +151,7 @@ impl PySubscriptionListener { if let Some(l) = guard.take() { l.close() .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; } Ok(()) }) @@ -192,10 +192,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; match l.next().await { Some(Ok(msg)) => Ok(PySubscriptionMessage::from(msg)), - Some(Err(e)) => Err(PyException::new_err(e.to_string())), + Some(Err(e)) => Err(crate::listen_err(e)), None => Err(PyStopAsyncIteration::new_err(())), } }) @@ -207,10 +207,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; match l.next().await { Some(Ok(msg)) => Ok(Some(PySubscriptionMessage::from(msg))), - Some(Err(e)) => Err(PyException::new_err(e.to_string())), + Some(Err(e)) => Err(crate::listen_err(e)), None => Ok(None), } }) @@ -222,10 +222,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.ack(&message_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(Python::attach(|py| py.None())) }) } @@ -236,10 +236,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.nack(&message_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(Python::attach(|py| py.None())) }) } @@ -254,10 +254,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.subscribe(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(Python::attach(|py| py.None())) }) } @@ -272,10 +272,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.unsubscribe(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(Python::attach(|py| py.None())) }) } @@ -290,10 +290,10 @@ impl PySubscriptionListenerAsync { let mut guard = listener.lock().await; let l = guard .as_mut() - .ok_or_else(|| PyException::new_err("listener is closed"))?; + .ok_or_else(|| PyValueError::new_err("listener is closed"))?; l.set_subscriptions(&external_ids) .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(Python::attach(|py| py.None())) }) } @@ -305,7 +305,7 @@ impl PySubscriptionListenerAsync { if let Some(l) = guard.take() { l.close() .await - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; } Ok(Python::attach(|py| py.None())) }) diff --git a/datahub_python_bindings/src/subscriptions/mod.rs b/datahub_python_bindings/src/subscriptions/mod.rs index aa775e0..e0a54a8 100644 --- a/datahub_python_bindings/src/subscriptions/mod.rs +++ b/datahub_python_bindings/src/subscriptions/mod.rs @@ -1,6 +1,7 @@ use crate::PyIdCollection; use crate::timeseries::PyTimeSeries; -use intellistream_datahub_sdk::filters::DataSort; +use crate::StringOrList; +use crate::events::PyTimeFilter; use intellistream_datahub_sdk::generic::IdAndExtId; use intellistream_datahub_sdk::subscriptions::{ DataCollectionString, DataWrapperMessage, EventAction, EventObject, Subscription, @@ -34,7 +35,9 @@ impl From for Subscription { } } -#[pyclass(module = "intellistream_datahub_sdk", name = "SubscriptionFilter")] +/// Criteria for `subscriptions.filter`. Every field is optional; the fields AND together and the +/// entries within a list OR. +#[pyclass(module = "intellistream_datahub_sdk", name = "SubscriptionFilter", from_py_object)] #[derive(Clone, Default)] pub struct PySubscriptionFilter { pub inner: SubscriptionFilter, @@ -53,19 +56,54 @@ impl From for SubscriptionFilter { #[pymethods] impl PySubscriptionFilter { + /// `external_id` and `name` are **pattern** lists: `*` and `%` are wildcards, `_` is literal, + /// matching is case-insensitive, and an entry with no wildcard matches exactly. Each also + /// accepts a bare string. `timeseries` matches subscriptions bound to at least one of them. #[new] - #[pyo3(signature=(timeseries=None))] - fn new(timeseries: Option>) -> Self { - let timeseries = timeseries - .unwrap_or_default() - .into_iter() - .map(IdAndExtId::from) - .collect(); + #[pyo3(signature = ( + id = None, + external_id = None, + name = None, + timeseries = None, + created_time = None, + last_updated_time = None, + ))] + pub fn new( + id: Option>, + external_id: Option, + name: Option, + timeseries: Option>, + created_time: Option, + last_updated_time: Option, + ) -> Self { Self { - inner: SubscriptionFilter { timeseries }, + inner: SubscriptionFilter { + id, + external_id: external_id.map(Into::into), + name: name.map(Into::into), + timeseries: timeseries + .unwrap_or_default() + .into_iter() + .map(IdAndExtId::from) + .collect(), + created_time: created_time.map(Into::into), + last_updated_time: last_updated_time.map(Into::into), + }, } } + #[getter] + fn id(&self) -> Option> { + self.inner.id.clone() + } + #[getter] + fn external_id(&self) -> Option> { + self.inner.external_id.clone() + } + #[getter] + fn name(&self) -> Option> { + self.inner.name.clone() + } #[getter] fn timeseries(&self) -> Vec { self.inner @@ -77,97 +115,43 @@ impl PySubscriptionFilter { } } -#[pyclass(module = "intellistream_datahub_sdk", name = "DataSort")] -#[derive(Clone, Default)] -pub struct PyDataSort { - pub inner: DataSort, -} - -impl From for PyDataSort { - fn from(s: DataSort) -> Self { - Self { inner: s } - } -} -impl From for DataSort { - fn from(s: PyDataSort) -> Self { - s.inner - } -} - -#[pymethods] -impl PyDataSort { - #[new] - #[pyo3(signature=(property=None, order=None))] - fn new(property: Option>, order: Option) -> Self { - Self { - inner: DataSort { - property: property.unwrap_or_default(), - order, - }, - } - } - - #[getter] - fn property(&self) -> Vec { - self.inner.property.clone() - } - #[getter] - fn order(&self) -> Option { - self.inner.order.clone() - } -} - -#[pyclass(module = "intellistream_datahub_sdk", name = "SubscriptionFilterForm")] -#[derive(Clone)] -pub struct PySubscriptionFilterForm { - pub inner: SubscriptionFilterForm, -} - -impl From for PySubscriptionFilterForm { - fn from(s: SubscriptionFilterForm) -> Self { - Self { inner: s } - } -} -impl From for SubscriptionFilterForm { - fn from(s: PySubscriptionFilterForm) -> Self { - s.inner - } -} - -#[pymethods] -impl PySubscriptionFilterForm { - #[new] - #[pyo3(signature=(filter=None, limit=None, sort=None))] - fn new( - filter: Option, - limit: Option, - sort: Option, - ) -> Self { - let mut inner = SubscriptionFilterForm::default(); - if let Some(f) = filter { - inner.filter = f.into(); - } - if let Some(l) = limit { - inner.limit = l; - } - if let Some(s) = sort { - inner.sort = Some(s.into()); - } - Self { inner } - } - - #[getter] - fn filter(&self) -> PySubscriptionFilter { - self.inner.filter.clone().into() - } - #[getter] - fn limit(&self) -> u32 { - self.inner.limit - } - #[getter] - fn sort(&self) -> Option { - self.inner.sort.clone().map(PyDataSort::from) - } +/// Build the request body for `subscriptions.filter` from either form of its arguments. +/// +/// Shared by the sync and async services so the accepted keywords cannot drift apart between them. +#[allow(clippy::too_many_arguments)] +pub fn subscription_filter_form( + filter: Option, + id: Option>, + external_id: Option, + name: Option, + timeseries: Option>, + created_time: Option, + last_updated_time: 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() + || timeseries.is_some() + || created_time.is_some() + || last_updated_time.is_some(); + let from_keywords = PySubscriptionFilter::new( + id, + external_id, + name, + timeseries, + created_time, + last_updated_time, + ) + .inner; + Ok(SubscriptionFilterForm { + filter: crate::resolve_filter(filter.map(Into::into), from_keywords, any_keyword)?, + limit, + paging: crate::build_page_request(sort_by, sort_order, cursor), + }) } /// Things accepted as a subscription identifier when deleting. @@ -463,8 +447,6 @@ impl PySubscriptionMessage { 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::()?; diff --git a/datahub_python_bindings/src/subscriptions/sync_service.rs b/datahub_python_bindings/src/subscriptions/sync_service.rs index 11c4536..c264200 100644 --- a/datahub_python_bindings/src/subscriptions/sync_service.rs +++ b/datahub_python_bindings/src/subscriptions/sync_service.rs @@ -1,14 +1,11 @@ use crate::subscriptions::listener::{PySubscriptionListener, shared_listener}; use crate::subscriptions::{ - PyDataSort, PySubscription, PySubscriptionFilterForm, SubscriptionIdentifyable, - SubscriptionTimeseriesId, + PySubscription, PySubscriptionFilter, SubscriptionIdentifyable, SubscriptionTimeseriesId, + subscription_filter_form, }; use intellistream_datahub_sdk::ApiService; use intellistream_datahub_sdk::generic::IdAndExtId; -use intellistream_datahub_sdk::subscriptions::{ - Subscription, SubscriptionFilter, SubscriptionFilterForm, -}; -use pyo3::exceptions::{PyException, PyValueError}; +use intellistream_datahub_sdk::subscriptions::Subscription; use pyo3::prelude::*; use std::sync::Arc; @@ -57,30 +54,42 @@ impl PySubscriptionsServiceSync { }) } - /// Subscriptions matching every criterion on the filter. - #[pyo3(signature=(form=None, *, timeseries=None, limit=None, sort=None))] + /// Subscriptions matching every criterion on the filter, newest first. + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, timeseries=None, + created_time=None, last_updated_time=None, limit=None, sort_by=None, + sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] fn filter( &self, py: Python<'_>, - form: Option, + filter: Option, + id: Option>, + external_id: Option, + name: Option, timeseries: Option>, - limit: Option, - sort: Option, - ) -> PyResult> { - let form = build_filter_form(form, timeseries, limit, sort)?; + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult { + let form = subscription_filter_form( + filter, id, external_id, name, timeseries, created_time, last_updated_time, limit, + sort_by, sort_order, cursor, + )?; let service = self.api_service.clone(); - py.detach(|| { + let (items, next_cursor) = py.detach(|| { let result = self .runtime .block_on(service.subscriptions.filter(&form)) - .map_err(|e| crate::datahub_err(e))?; - Ok(result - .get_items() - .iter() - .cloned() - .map(PySubscription::from) - .collect()) - }) + .map_err(crate::datahub_err)?; + let next_cursor = result.next_cursor().map(str::to_string); + let items: Vec = + result.get_items().iter().cloned().map(PySubscription::from).collect(); + Ok::<_, PyErr>((items, next_cursor)) + })?; + crate::PyPage::new(py, items, next_cursor) } fn delete(&self, py: Python<'_>, input: Vec) -> PyResult<()> { @@ -108,7 +117,7 @@ impl PySubscriptionsServiceSync { let listener = self .runtime .block_on(service.subscriptions.listen(&subscription_external_ids)) - .map_err(|e| PyException::new_err(e.to_string()))?; + .map_err(crate::listen_err)?; Ok(PySubscriptionListener { listener: shared_listener(listener), runtime, @@ -116,34 +125,3 @@ impl PySubscriptionsServiceSync { }) } } - -pub(crate) fn build_filter_form( - form: Option, - timeseries: Option>, - limit: Option, - sort: Option, -) -> PyResult { - let kwargs_used = timeseries.is_some() || limit.is_some() || sort.is_some(); - if form.is_some() && kwargs_used { - return Err(PyValueError::new_err( - "pass either a SubscriptionFilterForm or kwargs, not both", - )); - } - if let Some(r) = form { - return Ok(r.into()); - } - let mut r = SubscriptionFilterForm::default(); - if let Some(ts) = timeseries { - r.filter = SubscriptionFilter { - timeseries: ts.into_iter().map(IdAndExtId::from).collect(), - }; - } - if let Some(l) = limit { - r.limit = l; - } - if let Some(s) = sort { - r.sort = Some(s.into()); - } - Ok(r) -} - diff --git a/datahub_python_bindings/src/timeseries/async_service.rs b/datahub_python_bindings/src/timeseries/async_service.rs index 98bc811..54f50fa 100644 --- a/datahub_python_bindings/src/timeseries/async_service.rs +++ b/datahub_python_bindings/src/timeseries/async_service.rs @@ -2,12 +2,11 @@ use crate::timeseries::datapoints::{ PyDatapoint, PyDatapointsCollectionDatapoints, PyDatapointsCollectionString, }; use crate::timeseries::{ - PyDeleteFilter, PyTimeSeries, PyTimeSeriesUpdate, PyTimeseriesIdentifiable, + lists_to_collection, PyDeleteFilter, PyTimeSeries, PyTimeSeriesUpdate, PyTimeseriesIdentifiable, }; use crate::{ DatahubIdentity, Identifiable, PyIdCollection, PyRetrieveFilter, }; -use crate::datetime::py_datetime_to_utc; use intellistream_datahub_sdk::generic::{ DataWrapper, DatapointString, DatapointsCollection, DeleteFilter, IdAndExtId, RetrieveFilter, }; @@ -242,26 +241,9 @@ impl PyTimeSeriesServiceAsync { ts: Identifiable, ) -> PyResult> { let service = self.api_service.clone(); - let datapoints: Vec = timestamps - .into_iter() - .zip(values.into_iter()) - .map(|(timestamp, value)| { - Ok(DatapointString { - timestamp: py_datetime_to_utc(×tamp)?.timestamp_millis().to_string(), - value: value.to_string(), - }) - }) - .collect::>>()?; - let inner: DatapointsCollection = DatapointsCollection { - datapoints, - next_cursor: None, - id: ts.id_collection().id, - external_id: ts.id_collection().external_id, - unit: None, - unit_external_id: None, - }; + let collection = lists_to_collection(timestamps, values, ts)?; let mut wrapper = - DataWrapper::>::from_vec(vec![inner]); + DataWrapper::>::from_vec(vec![collection]); future_into_py(py, async move { let result = service .time_series diff --git a/datahub_python_bindings/src/timeseries/sync_service.rs b/datahub_python_bindings/src/timeseries/sync_service.rs index dad14d9..325aec5 100644 --- a/datahub_python_bindings/src/timeseries/sync_service.rs +++ b/datahub_python_bindings/src/timeseries/sync_service.rs @@ -1,5 +1,4 @@ use super::*; -use crate::datetime::py_datetime_to_utc; use crate::timeseries::datapoints::{ PyDatapointsCollectionDatapoints, PyDatapointsCollectionString, }; @@ -273,26 +272,9 @@ impl PyTimeSeriesServiceSync { ts: Identifiable, ) -> PyResult> { let service = self.api_service.clone(); - let datapoints: Vec = timestamps - .into_iter() - .zip(values.into_iter()) - .map(|(timestamp, value)| { - Ok(DatapointString { - timestamp: py_datetime_to_utc(×tamp)?.timestamp_millis().to_string(), - value: value.to_string(), - }) - }) - .collect::>>()?; - let inner: DatapointsCollection = DatapointsCollection { - datapoints, - next_cursor: None, - id: ts.id_collection().id, - external_id: ts.id_collection().external_id, - unit: None, - unit_external_id: None, - }; + let collection = lists_to_collection(timestamps, values, ts)?; let mut wrapper = - DataWrapper::>::from_vec(vec![inner]); + DataWrapper::>::from_vec(vec![collection]); py.detach(|| { let result = self .runtime diff --git a/python_tests/test_files.py b/python_tests/test_files.py index 05efc3a..3e2748c 100644 --- a/python_tests/test_files.py +++ b/python_tests/test_files.py @@ -159,6 +159,18 @@ def test_file_update_requires_a_selector(): intellistream_datahub_sdk.FileUpdate(name="renamed.jpg") +def test_file_upload_on_a_bad_path_raises_os_errors(tmp_path): + missing = str(tmp_path / "no_such_file.jpg") + with pytest.raises(FileNotFoundError): + intellistream_datahub_sdk.FileUpload(missing) + with pytest.raises(FileNotFoundError): + intellistream_datahub_sdk.FileUpload.from_path(missing) + with pytest.raises(FileNotFoundError): + intellistream_datahub_sdk.FileUpload.new_with_destination_path(missing, "/x") + with pytest.raises(IsADirectoryError): + intellistream_datahub_sdk.FileUpload(str(tmp_path)) + + # --------------------------------------------------------------------------- # # FileUpdate, field by field. # diff --git a/python_tests/test_subscriptions.py b/python_tests/test_subscriptions.py index af44a15..0278cca 100644 --- a/python_tests/test_subscriptions.py +++ b/python_tests/test_subscriptions.py @@ -51,13 +51,27 @@ def test_create_list_delete(sync_client, subscription_timeseries): filtered = sync_client.subscriptions.filter(timeseries=[ts_a_ext], limit=100) assert any(s.external_id == sub_ext for s in filtered) - # Same call via an explicit form. - form = intellistream_datahub_sdk.SubscriptionFilterForm( - filter=intellistream_datahub_sdk.SubscriptionFilter(timeseries=[ts_a_ext]), - limit=100, + # Same call via a prepared filter. + prepared = intellistream_datahub_sdk.SubscriptionFilter(timeseries=[ts_a_ext]) + filtered_via_filter = sync_client.subscriptions.filter(filter=prepared, limit=100) + assert any(s.external_id == sub_ext for s in filtered_via_filter) + + # The node-style criteria narrow to exactly this subscription. Name is case-insensitive. + for criteria in ( + {"id": [created[0].id]}, + {"external_id": sub_ext}, + {"name": f"*{sub_ext.upper()}"}, + ): + narrowed = sync_client.subscriptions.filter(**criteria) + assert [s.external_id for s in narrowed] == [sub_ext], criteria + + # A full page carries a cursor; continuing it under the same sort ends the walk. + first = sync_client.subscriptions.filter(external_id=sub_ext, limit=1, sort_by="externalId") + assert len(first) == 1 and first.next_cursor is not None + rest = sync_client.subscriptions.filter( + external_id=sub_ext, limit=1, sort_by="externalId", cursor=first.next_cursor ) - filtered_via_retriever = sync_client.subscriptions.filter(form) - assert any(s.external_id == sub_ext for s in filtered_via_retriever) + assert list(rest) == [] # Delete and verify gone. sync_client.subscriptions.delete([sub_ext]) @@ -86,10 +100,10 @@ def test_create_over_missing_timeseries_raises(sync_client): sync_client.subscriptions.create([sub]) -def test_filter_rejects_retriever_and_kwargs_together(sync_client): - form = intellistream_datahub_sdk.SubscriptionFilterForm() - with pytest.raises(ValueError): - sync_client.subscriptions.filter(form, limit=10) +def test_filter_rejects_filter_and_criteria_together(sync_client): + prepared = intellistream_datahub_sdk.SubscriptionFilter(name="anything") + with pytest.raises(TypeError): + sync_client.subscriptions.filter(filter=prepared, external_id="anything") def test_list_default_returns_list(sync_client): @@ -301,12 +315,13 @@ def test_listen_refused_subscription_surfaces_as_error(sync_client): bogus_sub = unique_id("sub_missing") listener = sync_client.subscriptions.listen([bogus_sub]) try: - with pytest.raises(Exception) as excinfo: + with pytest.raises(intellistream_datahub_sdk.DataHubException) as excinfo: # The server sends the error frame on attach, so the first iteration raises. for _ in listener: break message = str(excinfo.value) assert "not-found" in message or "forbidden" in message, message + assert excinfo.value.status_code is None finally: try: listener.close() diff --git a/python_tests/test_timeseries_datapoints.py b/python_tests/test_timeseries_datapoints.py index 637a160..e1c370c 100644 --- a/python_tests/test_timeseries_datapoints.py +++ b/python_tests/test_timeseries_datapoints.py @@ -357,3 +357,11 @@ def test_range_query_end_is_exclusive(sync_client, make_ts): assert not exclusive_end, ( f"end is exclusive, so a window closing on the point must not return it: {exclusive_end}" ) + + +def test_insert_from_lists_rejects_mismatched_lengths(sync_client, async_client): + at = datetime.datetime(2024, 1, 1, tzinfo=datetime.timezone.utc) + with pytest.raises(ValueError, match="got 2 and 1"): + sync_client.timeseries.insert_from_lists([at, at], [1.0], "never_sent") + with pytest.raises(ValueError, match="got 1 and 2"): + async_client.timeseries.insert_from_lists([at], [1.0, 2.0], "never_sent") diff --git a/src/datasets/mod.rs b/src/datasets/mod.rs index d844314..eacdbf9 100644 --- a/src/datasets/mod.rs +++ b/src/datasets/mod.rs @@ -73,7 +73,7 @@ impl DatasetsService { /// `POST /datasets/filter` — datasets matching [`DatasetFilterForm`], newest first. /// /// Every criterion on the filter is honoured server-side. Results are capped by the form's - /// `limit` (default 100 on the form, 1000 server-side when unset; max 10000). A broader match + /// `limit` (default 1000; max 10000). A broader match /// than the cap is paged, not truncated: set [`sort`](DatasetFilterForm) and echo the /// envelope's `next_cursor` back as `cursor` to walk it, under that same `sort`. pub async fn filter( @@ -503,7 +503,7 @@ impl DatasetFilterForm { pub fn new() -> Self { Self { filter: None, - limit: 100, + limit: 1000, paging: Default::default(), } } diff --git a/src/datasets/tests.rs b/src/datasets/tests.rs index 5cb7bfa..02c17fb 100644 --- a/src/datasets/tests.rs +++ b/src/datasets/tests.rs @@ -269,7 +269,7 @@ fn filter_body_matches_the_documented_wire_shape() { ); let json: serde_json::Value = serde_json::to_value(&filter).unwrap(); - assert_eq!(json["limit"], 100); + assert_eq!(json["limit"], 1000); let f = &json["filter"]; // The shared node criteria are flattened, so they sit directly on the filter body rather than // nested under a key the api does not read. diff --git a/src/files/mod.rs b/src/files/mod.rs index 2cbf2cc..810d0fa 100644 --- a/src/files/mod.rs +++ b/src/files/mod.rs @@ -9,6 +9,7 @@ use reqwest::Body; use serde::{Deserialize, Serialize}; use std::collections::HashMap; use std::fs; +use std::io; use std::path::Path; use std::sync::Weak; use tokio::fs::File; @@ -35,7 +36,12 @@ impl FileService { ) -> Result, ResponseError> { // The backend takes the file content as the raw PUT body; all metadata travels in // `X-Datahub-*` headers (see `FileController.upload`). - let body = file_upload.get_body().await; + let body = file_upload.get_body().await.map_err(|e| { + ResponseError::bad_request(format!( + "failed to open '{}': {}", + file_upload.file_path, e + )) + })?; let headers = file_upload.upload_headers(); self.execute_file_upload_request(self.base_url.as_str(), body, headers) .await @@ -364,19 +370,32 @@ pub struct FileUpload { } impl FileUpload { - pub fn new_with_destination_path(file_path: &str, destination_path: &str) -> Self { - let mut f = Self::new(file_path); + pub fn new_with_destination_path( + file_path: &str, + destination_path: &str, + ) -> io::Result { + let mut f = Self::new(file_path)?; f.set_destination_path(destination_path.to_string()); - f + Ok(f) } - pub fn new(file_path: &str) -> Self { - let metadata = fs::metadata(file_path).unwrap_or_else(|e| { - panic!("Failed to get metadata for file '{}': {}", file_path, e); - }); + /// Fails with the underlying `io::Error` when `file_path` cannot be read, and with + /// `ErrorKind::IsADirectory` or `InvalidInput` when it is not a regular file. + pub fn new(file_path: &str) -> io::Result { + let metadata = fs::metadata(file_path) + .map_err(|e| io::Error::new(e.kind(), format!("'{}': {}", file_path, e)))?; + if metadata.is_dir() { + return Err(io::Error::new( + io::ErrorKind::IsADirectory, + format!("'{}' is a directory, not a file", file_path), + )); + } if !metadata.is_file() { - panic!("Path '{}' is not a regular file.", file_path); + return Err(io::Error::new( + io::ErrorKind::InvalidInput, + format!("'{}' is not a regular file", file_path), + )); } let mut source_date_created: Option> = None; @@ -396,7 +415,12 @@ impl FileUpload { .file_name() .and_then(|name| name.to_str()) .map(|s| s.to_string()) - .unwrap_or_else(|| panic!("Could not get file name from path: {:?}", file_path)); + .ok_or_else(|| { + io::Error::new( + io::ErrorKind::InvalidInput, + format!("could not get a file name from '{}'", file_path), + ) + })?; let kind: Option = match infer::get_from_path(file_path) { Ok(Some(file_type)) => Some(file_type.mime_type().to_string()), @@ -410,7 +434,7 @@ impl FileUpload { } }; - Self { + Ok(Self { external_id: to_snake_lower_cased_allow_start_with_digits(file_name.as_str()), file_path: file_path.to_string(), destination_path: None, @@ -423,17 +447,15 @@ impl FileUpload { related_resources: None, source_date_created, source_last_updated, - } + }) } /// Opens the file and returns its contents as a streaming request body. The file content is /// the raw PUT body of the new `/files` upload endpoint. - pub async fn get_body(&self) -> Body { - let file = File::open(&self.file_path).await.unwrap_or_else(|e| { - panic!("Failed to open file '{}': {}", self.file_path, e); - }); + pub async fn get_body(&self) -> io::Result { + let file = File::open(&self.file_path).await?; let stream = FramedRead::new(file, BytesCodec::new()); - Body::wrap_stream(stream) + Ok(Body::wrap_stream(stream)) } /// Builds the `X-Datahub-*` and `Content-Type` headers the upload endpoint reads before it diff --git a/src/files/test.rs b/src/files/test.rs index 56d1e90..cbc18b2 100644 --- a/src/files/test.rs +++ b/src/files/test.rs @@ -6,6 +6,14 @@ mod tests { use crate::{create_api_service, ApiService}; + #[test] + fn file_upload_on_a_bad_path_is_an_error_not_a_panic() { + let missing = FileUpload::new("resources/test/no_such_file.jpg").unwrap_err(); + assert_eq!(missing.kind(), std::io::ErrorKind::NotFound); + let directory = FileUpload::new("resources/test").unwrap_err(); + assert_eq!(directory.kind(), std::io::ErrorKind::IsADirectory); + } + #[tokio::test] async fn test_file_upload() -> Result<(), Box> { let api_service = create_api_service(); @@ -16,25 +24,25 @@ mod tests { let mut upload_forms = vec![]; let file_path = "resources/test/random_values.csv"; - let file_upload_form = FileUpload::new_with_destination_path(file_path, "/foo/bar"); + let file_upload_form = FileUpload::new_with_destination_path(file_path, "/foo/bar").unwrap(); upload_forms.push(file_upload_form); let file_path = "resources/test/image.jpg"; - let mut file_upload_form = FileUpload::new_with_destination_path(file_path, "/images/"); + let mut file_upload_form = FileUpload::new_with_destination_path(file_path, "/images/").unwrap(); file_upload_form.set_file_name("sola.jpg".to_string()); file_upload_form.set_external_id("image_sola_jpg".to_string()); upload_forms.push(file_upload_form); let file_path = "resources/test/image2.jpg"; let mut file_upload_form = - FileUpload::new_with_destination_path(file_path, "/images/insects"); + FileUpload::new_with_destination_path(file_path, "/images/insects").unwrap(); file_upload_form.set_file_name("fly.jpg".to_string()); file_upload_form.set_external_id("image_fly_jpg".to_string()); upload_forms.push(file_upload_form); let file_path = "resources/test/image3.jpg"; let mut file_upload_form = - FileUpload::new_with_destination_path(file_path, "/images/norway/"); + FileUpload::new_with_destination_path(file_path, "/images/norway/").unwrap(); file_upload_form.set_file_name("teigland.jpg".to_string()); file_upload_form.set_external_id("image_teigland_bomlo_jpg".to_string()); upload_forms.push(file_upload_form); @@ -243,8 +251,8 @@ mod tests { let created_millis: i64 = 1_704_067_200_000; let updated_millis: i64 = 1_704_153_600_000; - let upload = FileUpload::new_with_destination_path("resources/test/image.jpg", "/dates"); - let body = upload.get_body().await; + let upload = FileUpload::new_with_destination_path("resources/test/image.jpg", "/dates").unwrap(); + let body = upload.get_body().await.unwrap(); let headers: Vec<(&str, String)> = vec![ ("X-Datahub-Path", "/dates/epoch.jpg".to_string()), ("X-Datahub-External-Id", ext_id.to_string()), @@ -388,7 +396,7 @@ mod tests { .await; let mut upload = - FileUpload::new_with_destination_path("resources/test/image.jpg", "/lifecycle"); + FileUpload::new_with_destination_path("resources/test/image.jpg", "/lifecycle").unwrap(); upload.set_file_name("sola.jpg".to_string()); upload.set_external_id(ext_id.to_string()); let source_bytes = std::fs::read("resources/test/image.jpg")?; diff --git a/src/filters.rs b/src/filters.rs index 9c57963..bf60796 100644 --- a/src/filters.rs +++ b/src/filters.rs @@ -278,6 +278,8 @@ pub enum TimeFilter { /// - **events** (`/events/filter`): `eventTime`, `createdTime`, `lastUpdatedTime`, `externalId`, /// `type`, `subType`, `status`, `source`, `dataSetId`. Default is `eventTime` **ascending** — the /// order the keyset pages in, so paging does not change it. +/// - **subscriptions** (`/subscriptions/filter`): `id`, `externalId`, `name`, `createdTime`, +/// `lastUpdatedTime`. Default is `createdTime` descending. /// /// An unsortable property falls back to the default rather than being rejected, so a misspelling /// returns the default order — visibly not what was asked for. Anything that is not exactly `desc` @@ -392,7 +394,7 @@ impl EventFilterForm { pub fn default() -> Self { Self { filter: None, - limit: 100, + limit: 1000, cursor: None, sort: None, advanced_filter: None, diff --git a/src/subscriptions/mod.rs b/src/subscriptions/mod.rs index d261c3d..6cf07d3 100644 --- a/src/subscriptions/mod.rs +++ b/src/subscriptions/mod.rs @@ -7,7 +7,7 @@ pub use listen::{ SubscriptionListener, SubscriptionMessage, WsDatapoint, }; -use crate::filters::DataSort; +use crate::filters::{PageRequest, TimeFilter}; use crate::generic::{ApiServiceProvider, DataHubEntity, DataWrapper, IdAndExtId}; use crate::http::ResponseError; use crate::ApiService; @@ -56,17 +56,12 @@ impl SubscriptionsService { .await } - /// `POST /subscriptions/filter` — subscriptions matching [`SubscriptionFilterForm`]. - /// - /// This was `POST /subscriptions/list`, whose body was subscriptions-only: a `limit` that - /// defaulted to 100 where the rest of the api defaulted to 1000, a sort property that reached - /// the query unvalidated, and no cursor, so a tenant past one page could not reach the rest. - /// The api moved it onto the family envelope and removed `/list` rather than aliasing it, so a - /// client that has not moved gets a 404. + /// `POST /subscriptions/filter` — subscriptions matching [`SubscriptionFilterForm`], newest + /// created first. /// - /// [`SubscriptionFilterForm`] is a strict subset of what the endpoint now accepts: it does not - /// yet carry the `cursor`, nor the `id`, `externalId`, `name`, `createdTime` and - /// `lastUpdatedTime` criteria the filter grew alongside `timeseries`. + /// Only subscriptions whose bound timeseries the caller can *all* read are returned. A match + /// broader than `limit` is paged, not truncated: echo the envelope's `next_cursor` back as + /// [`PageRequest::cursor`], under the same `sort`. pub async fn filter( &self, form: &SubscriptionFilterForm, @@ -138,29 +133,65 @@ impl DataHubEntity for Subscription { } } +/// Criteria for `POST /subscriptions/filter`, mirroring the api's `SubscriptionFilter`. +/// +/// Deliberately **not** a [`NodeFilter`](crate::filters::NodeFilter): a subscription is not a +/// node, and has no `source`, `labels` or `metadata` to match on. What it shares with the node +/// filters it shares by name and meaning — [`external_id`](Self::external_id) and +/// [`name`](Self::name) are case-insensitive pattern lists (`*` and `%` wildcards, `_` literal), +/// and the time windows are inclusive at both ends. +/// +/// Fields AND together, entries within a list OR, and `None` or an empty list places no +/// restriction. #[derive(Debug, Serialize, Deserialize, Clone, Default)] #[serde(rename_all = "camelCase")] pub struct SubscriptionFilter { + /// Max 1000. Sent as strings, like every id on the wire. + #[serde( + default, + skip_serializing_if = "Option::is_none", + with = "crate::serde_helper::opt_string_id_vec" + )] + pub id: Option>, + /// Max 1000. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub external_id: Option>, + /// Max 1000. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub name: Option>, + /// Subscriptions bound to at least one of these timeseries. Max 1000. #[serde(default, skip_serializing_if = "Vec::is_empty")] pub timeseries: Vec, + /// Matched against the subscription's `dateCreated`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub created_time: Option, + /// Matched against the subscription's `lastUpdated`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub last_updated_time: Option, } -#[derive(Debug, Serialize, Deserialize, Clone)] +/// Request body of `POST /subscriptions/filter`. +/// +/// Sortable by `id`, `externalId`, `name`, `createdTime` and `lastUpdatedTime`; the default is +/// `createdTime` descending. +#[derive(Debug, Serialize, Deserialize, Clone, Default)] #[serde(rename_all = "camelCase")] pub struct SubscriptionFilterForm { pub filter: SubscriptionFilter, - pub limit: u32, - /// Absent means the endpoint's default order, `createdTime` descending. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub sort: Option, + /// Defaults to 1000 server-side and is capped at 10000 — above that the request is rejected + /// with 400. + #[serde(skip_serializing_if = "Option::is_none")] + pub limit: Option, + /// Ordering and paging. Flattened, so `sort` and `cursor` sit beside `filter` and `limit`. + #[serde(flatten)] + pub paging: PageRequest, } -impl Default for SubscriptionFilterForm { - fn default() -> Self { - SubscriptionFilterForm { - filter: SubscriptionFilter::default(), - limit: 100, - sort: None, +impl SubscriptionFilterForm { + pub fn new(filter: SubscriptionFilter) -> Self { + Self { + filter, + ..Self::default() } } } diff --git a/src/subscriptions/test.rs b/src/subscriptions/test.rs index cdb049f..97972e4 100644 --- a/src/subscriptions/test.rs +++ b/src/subscriptions/test.rs @@ -7,7 +7,7 @@ mod tests { EventAction, EventObject, Subscription, SubscriptionFilter, SubscriptionFilterForm, }; - use crate::filters::DataSort; + use crate::filters::{PageRequest, TimeFilter}; use crate::timeseries::TimeSeries; use crate::{create_api_service, ApiService}; use reqwest::StatusCode; @@ -53,23 +53,43 @@ mod tests { #[test] fn test_filter_form_default_serializes_cleanly() { - // SubscriptionFilterForm::default() must serialize to a body the backend accepts: - // `{"filter":{},"limit":100}` — an unset sort is omitted rather than sent empty, and - // `nulls` is absent because the api removed the field and now 400s on it. - let json = serde_json::to_value(&SubscriptionFilterForm::default()).unwrap(); - assert_eq!(json["limit"], 100); - let filter_obj = json["filter"].as_object().unwrap(); - assert!(filter_obj.get("timeseries").is_none()); - assert!(json.get("sort").is_none()); - - let sorted = SubscriptionFilterForm { - sort: Some(DataSort::desc("name")), - ..Default::default() + // Every unset field is omitted rather than sent as null: the api rejects both unknown and + // null-typed keys. + let json = serde_json::to_value(SubscriptionFilterForm::default()).unwrap(); + assert_eq!(json, serde_json::json!({"filter": {}})); + } + + #[test] + fn test_filter_form_serializes_every_criterion_under_the_api_names() { + let min = "2026-01-01T00:00:00Z".parse().unwrap(); + let form = SubscriptionFilterForm { + filter: SubscriptionFilter { + id: Some(vec![12, 18]), + external_id: Some(vec!["plant_a_*".into()]), + name: Some(vec!["*dashboard*".into()]), + timeseries: vec![IdAndExtId::from_external_id("heater_2012_temp")], + created_time: Some(TimeFilter::After { min }), + last_updated_time: Some(TimeFilter::After { min }), + }, + limit: Some(100), + paging: PageRequest::asc("createdTime").after("opaque"), }; - let json = serde_json::to_value(&sorted).unwrap(); - assert_eq!(json["sort"]["property"], serde_json::json!(["name"])); - assert_eq!(json["sort"]["order"], "desc"); - assert!(json["sort"].as_object().unwrap().get("nulls").is_none()); + assert_eq!( + serde_json::to_value(&form).unwrap(), + serde_json::json!({ + "filter": { + "id": ["12", "18"], + "externalId": ["plant_a_*"], + "name": ["*dashboard*"], + "timeseries": [{"externalId": "heater_2012_temp"}], + "createdTime": {"min": "2026-01-01T00:00:00Z"}, + "lastUpdatedTime": {"min": "2026-01-01T00:00:00Z"}, + }, + "limit": 100, + "sort": {"property": ["createdTime"], "order": "asc"}, + "cursor": "opaque", + }) + ); } // Helpers for the integration test — mirrors delete_events in events/tests.rs. @@ -158,9 +178,10 @@ mod tests { .filter(&SubscriptionFilterForm { filter: SubscriptionFilter { timeseries: vec![IdAndExtId::from_external_id(&ts_a_ext)], + ..Default::default() }, - limit: 100, - sort: None, + limit: Some(100), + ..Default::default() }) .await?; assert!( @@ -168,6 +189,61 @@ mod tests { "timeseries-filtered list must contain the subscription" ); + // 4b. The node-style criteria narrow too. The name is matched case-insensitively. + let created_id = created_item.id.unwrap(); + let created_at = created_item.date_created.unwrap(); + let narrow = |filter: SubscriptionFilter| SubscriptionFilterForm { + filter, + ..Default::default() + }; + for (label, filter) in [ + ("id", SubscriptionFilter { id: Some(vec![created_id]), ..Default::default() }), + ( + "externalId", + SubscriptionFilter { external_id: Some(vec![sub_ext.clone()]), ..Default::default() }, + ), + ( + "name", + SubscriptionFilter { + name: Some(vec![format!("*{}", sub_ext.to_uppercase())]), + ..Default::default() + }, + ), + ] { + let found = api_service.subscriptions.filter(&narrow(filter)).await?; + let ext_ids: Vec<&str> = + found.get_items().iter().map(|s| s.external_id.as_str()).collect(); + assert_eq!(ext_ids, vec![sub_ext.as_str()], "filter by {label}"); + } + let too_late = api_service + .subscriptions + .filter(&narrow(SubscriptionFilter { + external_id: Some(vec![sub_ext.clone()]), + created_time: Some(TimeFilter::After { + min: created_at + chrono::Duration::hours(1), + }), + ..Default::default() + })) + .await?; + assert!(too_late.get_items().is_empty(), "createdTime must narrow"); + + // 4c. A full page carries a cursor; continuing it under the same sort ends the walk. + let mut paged = narrow(SubscriptionFilter { + external_id: Some(vec![sub_ext.clone()]), + ..Default::default() + }); + paged.limit = Some(1); + paged.paging = PageRequest::asc("externalId"); + let first = api_service.subscriptions.filter(&paged).await?; + assert_eq!(first.length(), 1); + let cursor = first + .next_cursor() + .expect("a full page must carry a next cursor") + .to_string(); + paged.paging = PageRequest::asc("externalId").after(&cursor); + let second = api_service.subscriptions.filter(&paged).await?; + assert!(second.get_items().is_empty(), "the walk ends after the only match"); + // 5. Delete the subscription. Backend returns 204 No Content → empty DataWrapper. delete_subscriptions(&api_service, &[IdAndExtId::from_external_id(&sub_ext)]).await; // Explicit delete succeeded — disarm the guard so it doesn't re-delete. @@ -179,9 +255,10 @@ mod tests { .filter(&SubscriptionFilterForm { filter: SubscriptionFilter { timeseries: vec![IdAndExtId::from_external_id(&ts_a_ext)], + ..Default::default() }, - limit: 100, - sort: None, + limit: Some(100), + ..Default::default() }) .await?; assert!(