From ddbe033817b0d67f3391a91053e6ddf22f4cdb02 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 10:57:27 +0200 Subject: [PATCH 1/7] fix: datasets, events and subscriptions filters default to 1000 DatasetFilterForm, EventFilterForm and SubscriptionFilterForm all sent limit=100 when the caller left it unset, and the Python events.filter did the same through unwrap_or(100). timeseries.filter and every list() use the server's 1000, so the page size depended on which entity was asked about. All three now default to 1000. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- datahub_python_bindings/src/datasets/mod.rs | 2 +- datahub_python_bindings/src/events/mod.rs | 4 +++- src/datasets/mod.rs | 4 ++-- src/datasets/tests.rs | 2 +- src/filters.rs | 2 +- src/subscriptions/mod.rs | 2 +- src/subscriptions/test.rs | 4 ++-- 7 files changed, 11 insertions(+), 9 deletions(-) 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/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/filters.rs b/src/filters.rs index 9c57963..8debaa5 100644 --- a/src/filters.rs +++ b/src/filters.rs @@ -392,7 +392,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..4dd0f50 100644 --- a/src/subscriptions/mod.rs +++ b/src/subscriptions/mod.rs @@ -159,7 +159,7 @@ impl Default for SubscriptionFilterForm { fn default() -> Self { SubscriptionFilterForm { filter: SubscriptionFilter::default(), - limit: 100, + limit: 1000, sort: None, } } diff --git a/src/subscriptions/test.rs b/src/subscriptions/test.rs index cdb049f..25ffb49 100644 --- a/src/subscriptions/test.rs +++ b/src/subscriptions/test.rs @@ -54,10 +54,10 @@ 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 + // `{"filter":{},"limit":1000}` — 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); + assert_eq!(json["limit"], 1000); let filter_obj = json["filter"].as_object().unwrap(); assert!(filter_obj.get("timeseries").is_none()); assert!(json.get("sort").is_none()); From 6268fe5cf194256098b5b6bf06872a4a7a9aae20 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 10:58:30 +0200 Subject: [PATCH 2/7] feat!: files.search takes a limit GET /files/search reads a limit parameter (default 100, clamped to 1000), but FileService::search never sent one, so a search could not return more than 100 hits. It now takes `limit: Option`, omitted when None; Python's files.search gains `limit=None` on both clients. BREAKING CHANGE: FileService::search (and the blocking mirror) takes a second argument; pass None for the previous behaviour. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- .../python/intellistream_datahub_sdk/__init__.pyi | 4 ++-- .../src/files/async_service.rs | 13 ++++++++++--- datahub_python_bindings/src/files/sync_service.rs | 8 +++++--- python_tests/test_files.py | 1 + src/blocking.rs | 2 +- src/files/mod.rs | 15 ++++++++++++--- src/files/test.rs | 6 ++++-- 7 files changed, 35 insertions(+), 14 deletions(-) diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index ccde6d0..3c8bef8 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -2036,7 +2036,7 @@ class FilesServiceSync: def list_directory_by_path(self, path: str) -> list[INode]: ... def get_by_id(self, id: int) -> list[INode]: ... def get_by_external_id(self, external_id: str) -> list[INode]: ... - def search(self, query: str) -> list[INode]: ... + def search(self, query: str, limit: int | None = None) -> list[INode]: ... def list_trash(self) -> list[INode]: ... def restore(self, input: list[FileIdentifiable]) -> list[INode]: ... def update(self, update: FileUpdate) -> list[INode]: ... @@ -2051,7 +2051,7 @@ class FilesServiceAsync: async def list_directory_by_path(self, path: str) -> list[INode]: ... async def get_by_id(self, id: int) -> list[INode]: ... async def get_by_external_id(self, external_id: str) -> list[INode]: ... - async def search(self, query: str) -> list[INode]: ... + async def search(self, query: str, limit: int | None = None) -> list[INode]: ... async def list_trash(self) -> list[INode]: ... async def restore(self, input: list[FileIdentifiable]) -> list[INode]: ... async def update(self, update: FileUpdate) -> list[INode]: ... diff --git a/datahub_python_bindings/src/files/async_service.rs b/datahub_python_bindings/src/files/async_service.rs index 06c884c..b97570c 100644 --- a/datahub_python_bindings/src/files/async_service.rs +++ b/datahub_python_bindings/src/files/async_service.rs @@ -139,13 +139,20 @@ impl PyFilesServiceAsync { }) } - /// Full-text search over file and folder names and descriptions. - fn search<'py>(&self, py: Python<'py>, query: String) -> PyResult> { + /// Full-text search over file and folder names and descriptions. `limit` defaults to 100 + /// server-side and is clamped to 1000. + #[pyo3(signature = (query, limit = None))] + fn search<'py>( + &self, + py: Python<'py>, + query: String, + limit: Option, + ) -> PyResult> { let service = self.api_service.clone(); future_into_py(py, async move { let result = service .files - .search(query.as_str()) + .search(query.as_str(), limit) .await .map_err(|e| crate::datahub_err(e))?; Ok(to_py_inodes(&result, &service)) diff --git a/datahub_python_bindings/src/files/sync_service.rs b/datahub_python_bindings/src/files/sync_service.rs index d4bdaf3..76352e0 100644 --- a/datahub_python_bindings/src/files/sync_service.rs +++ b/datahub_python_bindings/src/files/sync_service.rs @@ -127,13 +127,15 @@ impl PyFilesServiceSync { }) } - /// Full-text search over file and folder names and descriptions. - fn search<'py>(&self, py: Python<'py>, query: &str) -> PyResult> { + /// Full-text search over file and folder names and descriptions. `limit` defaults to 100 + /// server-side and is clamped to 1000. + #[pyo3(signature = (query, limit = None))] + fn search<'py>(&self, py: Python<'py>, query: &str, limit: Option) -> PyResult> { let service = self.api_service.clone(); py.detach(|| { let result = self .runtime - .block_on(service.files.search(query)) + .block_on(service.files.search(query, limit)) .map_err(|e| crate::datahub_err(e))?; Ok(to_py_inodes(&result, &service)) }) diff --git a/python_tests/test_files.py b/python_tests/test_files.py index 05efc3a..db09fbb 100644 --- a/python_tests/test_files.py +++ b/python_tests/test_files.py @@ -108,6 +108,7 @@ def test_get_search_update_download_trash_restore(sync_client, tmp_path): found = sync_client.files.search("sola") assert any(node.external_id == ext_id for node in found) + assert len(sync_client.files.search("sola", limit=1)) == 1 # A blank query is answered with an empty list, not an error. assert sync_client.files.search("") == [] diff --git a/src/blocking.rs b/src/blocking.rs index aec476c..e1c54cc 100644 --- a/src/blocking.rs +++ b/src/blocking.rs @@ -319,7 +319,7 @@ impl FileService { fn delete(id_collection: &DataWrapper) -> Result, ResponseError>; fn get_by_id(id: u64) -> Result, ResponseError>; fn get_by_external_id(external_id: &str) -> Result, ResponseError>; - fn search(query: &str) -> Result, ResponseError>; + fn search(query: &str, limit: Option) -> Result, ResponseError>; fn list_trash() -> Result, ResponseError>; fn restore(id_collection: &DataWrapper) -> Result, ResponseError>; fn update(update: &FileUpdate) -> Result, ResponseError>; diff --git a/src/files/mod.rs b/src/files/mod.rs index 2cbf2cc..ee3b186 100644 --- a/src/files/mod.rs +++ b/src/files/mod.rs @@ -92,10 +92,19 @@ impl FileService { /// `GET /files/search?q=` — case-insensitive full-text search over file and folder names and /// descriptions, across the whole tree, narrowed to the caller's readable datasets. /// - /// A blank query is answered with an empty item list rather than an error. - pub async fn search(&self, query: &str) -> Result, ResponseError> { + /// A blank query is answered with an empty item list rather than an error. `limit` defaults to + /// 100 server-side and is clamped to 1000 rather than rejected above it; `<= 0` is the default. + pub async fn search( + &self, + query: &str, + limit: Option, + ) -> Result, ResponseError> { let full_path = format!("{}/search", self.base_url.as_str()); - self.execute_get_request(full_path.as_str(), Some(&[("q", query.to_string())])) + let mut params = vec![("q", query.to_string())]; + if let Some(limit) = limit { + params.push(("limit", limit.to_string())); + } + self.execute_get_request(full_path.as_str(), Some(¶ms)) .await } diff --git a/src/files/test.rs b/src/files/test.rs index 56d1e90..3750b20 100644 --- a/src/files/test.rs +++ b/src/files/test.rs @@ -409,15 +409,17 @@ mod tests { assert_eq!(by_ext.get_items()[0].id, Some(id)); // --- search --- - let found = api_service.files.search("sola").await?; + let found = api_service.files.search("sola", None).await?; assert_eq!(found.get_http_status_code().unwrap(), 200); assert!( found.get_items().iter().any(|n| n.external_id == ext_id), "search for 'sola' should surface the uploaded file" ); + let capped = api_service.files.search("sola", Some(1)).await?; + assert_eq!(capped.get_items().len(), 1); // A blank query is answered with an empty list, not an error. - let empty = api_service.files.search("").await?; + let empty = api_service.files.search("", None).await?; assert_eq!(empty.get_http_status_code().unwrap(), 200); // --- download, in memory and streamed to disk --- From 2d1033f38b976a13368f916857dc4699b48d6dbe Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 10:59:52 +0200 Subject: [PATCH 3/7] fix(python): insert_from_lists rejects lists of different lengths The sync and async insert_from_lists zipped timestamps with values, so extra entries on the longer side were dropped without a word. Both now build their collection through lists_to_collection, as the binary twin already did, which raises ValueError naming the two lengths. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- .../src/timeseries/async_service.rs | 24 +++---------------- .../src/timeseries/sync_service.rs | 22 ++--------------- python_tests/test_timeseries_datapoints.py | 8 +++++++ 3 files changed, 13 insertions(+), 41 deletions(-) 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_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") From 03efe8980496d623f00e6271bbaa36735b02a9f4 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 11:01:39 +0200 Subject: [PATCH 4/7] fix(python): subscription errors raise the types the other services do subscriptions.filter raised ValueError when handed both a form and keywords, where every other filter raises TypeError; it now raises TypeError too. The listener raised a bare Exception for every WebSocket failure, so `except DataHubException` never caught one. listen() and every listener method now raise DataHubException, with status_code and the problem attributes set to None since there is no HTTP response behind them. Using a listener after close() raises ValueError, as a closed Python file does. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- .../intellistream_datahub_sdk/__init__.pyi | 6 +- datahub_python_bindings/src/lib.rs | 16 +++++ .../src/subscriptions/async_service.rs | 3 +- .../src/subscriptions/listener.rs | 62 +++++++++---------- .../src/subscriptions/sync_service.rs | 6 +- python_tests/test_subscriptions.py | 5 +- 6 files changed, 59 insertions(+), 39 deletions(-) diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index 3c8bef8..d473fb5 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 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..9239646 100644 --- a/datahub_python_bindings/src/subscriptions/async_service.rs +++ b/datahub_python_bindings/src/subscriptions/async_service.rs @@ -7,7 +7,6 @@ use crate::subscriptions::{ 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; @@ -117,7 +116,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/sync_service.rs b/datahub_python_bindings/src/subscriptions/sync_service.rs index 11c4536..10b0f5f 100644 --- a/datahub_python_bindings/src/subscriptions/sync_service.rs +++ b/datahub_python_bindings/src/subscriptions/sync_service.rs @@ -8,7 +8,7 @@ use intellistream_datahub_sdk::generic::IdAndExtId; use intellistream_datahub_sdk::subscriptions::{ Subscription, SubscriptionFilter, SubscriptionFilterForm, }; -use pyo3::exceptions::{PyException, PyValueError}; +use pyo3::exceptions::PyTypeError; use pyo3::prelude::*; use std::sync::Arc; @@ -108,7 +108,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, @@ -125,7 +125,7 @@ pub(crate) fn build_filter_form( ) -> PyResult { let kwargs_used = timeseries.is_some() || limit.is_some() || sort.is_some(); if form.is_some() && kwargs_used { - return Err(PyValueError::new_err( + return Err(PyTypeError::new_err( "pass either a SubscriptionFilterForm or kwargs, not both", )); } diff --git a/python_tests/test_subscriptions.py b/python_tests/test_subscriptions.py index af44a15..78cf134 100644 --- a/python_tests/test_subscriptions.py +++ b/python_tests/test_subscriptions.py @@ -88,7 +88,7 @@ def test_create_over_missing_timeseries_raises(sync_client): def test_filter_rejects_retriever_and_kwargs_together(sync_client): form = intellistream_datahub_sdk.SubscriptionFilterForm() - with pytest.raises(ValueError): + with pytest.raises(TypeError): sync_client.subscriptions.filter(form, limit=10) @@ -301,12 +301,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() From 797b2ae28f2774046e8c4b925123a233bda15d29 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 11:03:18 +0200 Subject: [PATCH 5/7] fix!: FileUpload on a bad path is an error, not a panic FileUpload::new panicked when the path was missing, a directory or had no file name, and get_body panicked when the file could not be opened. From Python that surfaced as PanicException, which derives from BaseException, so `except Exception` did not catch it. new and new_with_destination_path now return io::Result, and get_body io::Result; upload_file reports an open failure as a 400 ResponseError. In Python the constructors raise FileNotFoundError, IsADirectoryError or another OSError through PyO3's io::Error mapping. BREAKING CHANGE: FileUpload::new, FileUpload::new_with_destination_path and FileUpload::get_body return io::Result. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- .../intellistream_datahub_sdk/__init__.pyi | 3 + datahub_python_bindings/src/files/mod.rs | 6 +- python_tests/test_files.py | 12 ++++ src/files/mod.rs | 56 +++++++++++++------ src/files/test.rs | 22 +++++--- 5 files changed, 72 insertions(+), 27 deletions(-) diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index d473fb5..5af63ac 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -1951,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, 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/python_tests/test_files.py b/python_tests/test_files.py index db09fbb..4849666 100644 --- a/python_tests/test_files.py +++ b/python_tests/test_files.py @@ -160,6 +160,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/src/files/mod.rs b/src/files/mod.rs index ee3b186..89323b0 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 @@ -373,19 +379,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; @@ -405,7 +424,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()), @@ -419,7 +443,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, @@ -432,17 +456,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 3750b20..8f59f99 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")?; From df3ce1a7e1fb50a1165db8e864ffca2a22bf6852 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 13:08:11 +0200 Subject: [PATCH 6/7] Revert "feat!: files.search takes a limit" This reverts commit 6268fe5. files.search is left as it was. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- .../python/intellistream_datahub_sdk/__init__.pyi | 4 ++-- .../src/files/async_service.rs | 13 +++---------- datahub_python_bindings/src/files/sync_service.rs | 8 +++----- python_tests/test_files.py | 1 - src/blocking.rs | 2 +- src/files/mod.rs | 15 +++------------ src/files/test.rs | 6 ++---- 7 files changed, 14 insertions(+), 35 deletions(-) diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index 5af63ac..22161e4 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -2043,7 +2043,7 @@ class FilesServiceSync: def list_directory_by_path(self, path: str) -> list[INode]: ... def get_by_id(self, id: int) -> list[INode]: ... def get_by_external_id(self, external_id: str) -> list[INode]: ... - def search(self, query: str, limit: int | None = None) -> list[INode]: ... + def search(self, query: str) -> list[INode]: ... def list_trash(self) -> list[INode]: ... def restore(self, input: list[FileIdentifiable]) -> list[INode]: ... def update(self, update: FileUpdate) -> list[INode]: ... @@ -2058,7 +2058,7 @@ class FilesServiceAsync: async def list_directory_by_path(self, path: str) -> list[INode]: ... async def get_by_id(self, id: int) -> list[INode]: ... async def get_by_external_id(self, external_id: str) -> list[INode]: ... - async def search(self, query: str, limit: int | None = None) -> list[INode]: ... + async def search(self, query: str) -> list[INode]: ... async def list_trash(self) -> list[INode]: ... async def restore(self, input: list[FileIdentifiable]) -> list[INode]: ... async def update(self, update: FileUpdate) -> list[INode]: ... diff --git a/datahub_python_bindings/src/files/async_service.rs b/datahub_python_bindings/src/files/async_service.rs index b97570c..06c884c 100644 --- a/datahub_python_bindings/src/files/async_service.rs +++ b/datahub_python_bindings/src/files/async_service.rs @@ -139,20 +139,13 @@ impl PyFilesServiceAsync { }) } - /// Full-text search over file and folder names and descriptions. `limit` defaults to 100 - /// server-side and is clamped to 1000. - #[pyo3(signature = (query, limit = None))] - fn search<'py>( - &self, - py: Python<'py>, - query: String, - limit: Option, - ) -> PyResult> { + /// Full-text search over file and folder names and descriptions. + fn search<'py>(&self, py: Python<'py>, query: String) -> PyResult> { let service = self.api_service.clone(); future_into_py(py, async move { let result = service .files - .search(query.as_str(), limit) + .search(query.as_str()) .await .map_err(|e| crate::datahub_err(e))?; Ok(to_py_inodes(&result, &service)) diff --git a/datahub_python_bindings/src/files/sync_service.rs b/datahub_python_bindings/src/files/sync_service.rs index 76352e0..d4bdaf3 100644 --- a/datahub_python_bindings/src/files/sync_service.rs +++ b/datahub_python_bindings/src/files/sync_service.rs @@ -127,15 +127,13 @@ impl PyFilesServiceSync { }) } - /// Full-text search over file and folder names and descriptions. `limit` defaults to 100 - /// server-side and is clamped to 1000. - #[pyo3(signature = (query, limit = None))] - fn search<'py>(&self, py: Python<'py>, query: &str, limit: Option) -> PyResult> { + /// Full-text search over file and folder names and descriptions. + fn search<'py>(&self, py: Python<'py>, query: &str) -> PyResult> { let service = self.api_service.clone(); py.detach(|| { let result = self .runtime - .block_on(service.files.search(query, limit)) + .block_on(service.files.search(query)) .map_err(|e| crate::datahub_err(e))?; Ok(to_py_inodes(&result, &service)) }) diff --git a/python_tests/test_files.py b/python_tests/test_files.py index 4849666..3e2748c 100644 --- a/python_tests/test_files.py +++ b/python_tests/test_files.py @@ -108,7 +108,6 @@ def test_get_search_update_download_trash_restore(sync_client, tmp_path): found = sync_client.files.search("sola") assert any(node.external_id == ext_id for node in found) - assert len(sync_client.files.search("sola", limit=1)) == 1 # A blank query is answered with an empty list, not an error. assert sync_client.files.search("") == [] diff --git a/src/blocking.rs b/src/blocking.rs index e1c54cc..aec476c 100644 --- a/src/blocking.rs +++ b/src/blocking.rs @@ -319,7 +319,7 @@ impl FileService { fn delete(id_collection: &DataWrapper) -> Result, ResponseError>; fn get_by_id(id: u64) -> Result, ResponseError>; fn get_by_external_id(external_id: &str) -> Result, ResponseError>; - fn search(query: &str, limit: Option) -> Result, ResponseError>; + fn search(query: &str) -> Result, ResponseError>; fn list_trash() -> Result, ResponseError>; fn restore(id_collection: &DataWrapper) -> Result, ResponseError>; fn update(update: &FileUpdate) -> Result, ResponseError>; diff --git a/src/files/mod.rs b/src/files/mod.rs index 89323b0..810d0fa 100644 --- a/src/files/mod.rs +++ b/src/files/mod.rs @@ -98,19 +98,10 @@ impl FileService { /// `GET /files/search?q=` — case-insensitive full-text search over file and folder names and /// descriptions, across the whole tree, narrowed to the caller's readable datasets. /// - /// A blank query is answered with an empty item list rather than an error. `limit` defaults to - /// 100 server-side and is clamped to 1000 rather than rejected above it; `<= 0` is the default. - pub async fn search( - &self, - query: &str, - limit: Option, - ) -> Result, ResponseError> { + /// A blank query is answered with an empty item list rather than an error. + pub async fn search(&self, query: &str) -> Result, ResponseError> { let full_path = format!("{}/search", self.base_url.as_str()); - let mut params = vec![("q", query.to_string())]; - if let Some(limit) = limit { - params.push(("limit", limit.to_string())); - } - self.execute_get_request(full_path.as_str(), Some(¶ms)) + self.execute_get_request(full_path.as_str(), Some(&[("q", query.to_string())])) .await } diff --git a/src/files/test.rs b/src/files/test.rs index 8f59f99..cbc18b2 100644 --- a/src/files/test.rs +++ b/src/files/test.rs @@ -417,17 +417,15 @@ mod tests { assert_eq!(by_ext.get_items()[0].id, Some(id)); // --- search --- - let found = api_service.files.search("sola", None).await?; + let found = api_service.files.search("sola").await?; assert_eq!(found.get_http_status_code().unwrap(), 200); assert!( found.get_items().iter().any(|n| n.external_id == ext_id), "search for 'sola' should surface the uploaded file" ); - let capped = api_service.files.search("sola", Some(1)).await?; - assert_eq!(capped.get_items().len(), 1); // A blank query is answered with an empty list, not an error. - let empty = api_service.files.search("", None).await?; + let empty = api_service.files.search("").await?; assert_eq!(empty.get_http_status_code().unwrap(), 200); // --- download, in memory and streamed to disk --- From 846b470726df9f21cd6c74f44e006323229914a6 Mon Sep 17 00:00:00 2001 From: jgjesdal Date: Fri, 25 Sep 2026 14:38:51 +0200 Subject: [PATCH 7/7] feat!: subscriptions.filter takes the api's full criteria and pages SubscriptionFilter carries id, externalId, name, timeseries, createdTime and lastUpdatedTime, and SubscriptionFilterForm takes the family shape: filter, limit: Option, and a flattened PageRequest for sort and cursor. In Python, filter() takes the criteria as keywords or filter= plus sort_by/sort_order/cursor and returns a Page, like datasets.filter. SubscriptionFilterForm and DataSort are removed. Co-Authored-By: Claude Opus 5.5 (1M context) Signed-off-by: jgjesdal --- AGENTS.md | 2 +- .../intellistream_datahub_sdk/__init__.pyi | 73 ++++--- .../src/subscriptions/async_service.rs | 42 ++-- .../src/subscriptions/mod.rs | 188 ++++++++---------- .../src/subscriptions/sync_service.rs | 84 +++----- python_tests/test_subscriptions.py | 32 ++- src/filters.rs | 2 + src/subscriptions/mod.rs | 75 +++++-- src/subscriptions/test.rs | 119 +++++++++-- 9 files changed, 364 insertions(+), 253 deletions(-) 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 22161e4..c88c205 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -2082,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: @@ -2203,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: ... @@ -2217,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/subscriptions/async_service.rs b/datahub_python_bindings/src/subscriptions/async_service.rs index 9239646..c4b7570 100644 --- a/datahub_python_bindings/src/subscriptions/async_service.rs +++ b/datahub_python_bindings/src/subscriptions/async_service.rs @@ -1,8 +1,7 @@ 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; @@ -61,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)) }) } 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 10b0f5f..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::PyTypeError; +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<()> { @@ -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(PyTypeError::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/python_tests/test_subscriptions.py b/python_tests/test_subscriptions.py index 78cf134..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() +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(form, limit=10) + sync_client.subscriptions.filter(filter=prepared, external_id="anything") def test_list_default_returns_list(sync_client): diff --git a/src/filters.rs b/src/filters.rs index 8debaa5..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` diff --git a/src/subscriptions/mod.rs b/src/subscriptions/mod.rs index 4dd0f50..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: 1000, - 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 25ffb49..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":1000}` — 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"], 1000); - 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!(