diff --git a/AGENTS.md b/AGENTS.md index 4c62676..adf33e1 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -105,6 +105,15 @@ and `TimeSeriesFilterForm` — four spellings of one idea, which is what made `B first place. The backend still calls the request bodies `XRetreiver` (its own spelling of *Retriever*); the wire format is unaffected either way, since none of these names is serialized. +**That is the Rust surface. Python does not mirror it, deliberately.** There is no `XFilterForm` +in the bindings: `filter()` takes either the criteria as keywords — `client.timeseries.filter( +name=["Pump*"], limit=100)` — or a prepared `filter=` object, and passing both is a `TypeError`. +`limit`, `sort_by`, `sort_order` and `cursor` are always arguments of the call, never fields of the +filter, so one `XFilter` can be reused across `filter()` and `search()` without carrying paging +state between them. Before 0.3.0 the Python side had three shapes at once — an envelope for events +and datasets, a flat `TimeSeriesFilterForm` that was really the criteria, and bare keywords on +resources — and `timeseries.search` silently discarded the paging fields of the form it was handed. + The four `/{entity}/filter` endpoints share one contract. `NodeFilter` (`src/filters.rs`) is the criteria every node type can be filtered by — `id`, `externalId`, `name`, `source`, `labels`, `metadata`, `createdTime`, `lastUpdatedTime` — and `ResourceFilter`, `TimeSeriesFilter` and `DatasetFilter` each `#[serde(flatten)]` it, so on the wire its fields sit alongside the type-specific ones. `EventFilter` deliberately does **not** extend it (events are not nodes: no `name` column, a UUID id) but matches it field for field wherever ClickHouse can back it. The rules, which every one of them obeys: diff --git a/datahub_python_bindings/README.md b/datahub_python_bindings/README.md index 8b5b15e..d73d4f0 100644 --- a/datahub_python_bindings/README.md +++ b/datahub_python_bindings/README.md @@ -23,11 +23,18 @@ client.timeseries.insert_from_lists( ) # Find the alarms. -alarms = client.events.filter( - dh.EventFilterForm(filter=dh.EventFilter(type="ALARM"), limit=100) -) +alarms = client.events.filter(type="ALARM", limit=100) + +# The same criteria, kept as an object, so one definition can be filtered and searched with. +ALARMS = dh.EventFilter(type="ALARM") +recent = client.events.filter(filter=ALARMS, limit=100, sort_by="eventTime", sort_order="desc") +matches = client.events.search("bearing", filter=ALARMS) ``` +Every `filter()` takes either the criteria as keywords or a prepared `filter=` object — passing +both is a `TypeError`. `limit`, `sort_by`, `sort_order` and `cursor` are always arguments of the +call rather than fields of the filter, so a stored filter carries no paging state into its next use. + Both a synchronous and an asynchronous client are available — `DataHubClient` and `AsyncDataHubClient`. The async one exposes the same services with awaitable methods. diff --git a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi index 6c8baf2..509c2b3 100644 --- a/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi +++ b/datahub_python_bindings/python/intellistream_datahub_sdk/__init__.pyi @@ -42,14 +42,14 @@ class Page(Sequence[Any]): The whole loop is:: - page = client.timeseries.filter(TimeSeriesFilterForm(limit=100, sort_by="name")) + page = client.timeseries.filter(limit=100, sort_by="name") while True: for ts in page: ... if page.next_cursor is None: break - page = client.timeseries.filter(TimeSeriesFilterForm( - limit=100, sort_by="name", cursor=page.next_cursor)) + page = client.timeseries.filter( + limit=100, sort_by="name", cursor=page.next_cursor) ``next_cursor`` is ``None`` on the last page. A *full* page may still be the last one — the server does not count the rows twice — so a walk ends with one request that comes back empty. @@ -223,10 +223,14 @@ class IdCollection: def external_id(self) -> str | None: ... -class TimeSeriesFilterForm: +class TimeSeriesFilter: """AND-combined criteria for ``timeseries.filter`` (``POST /timeseries/filter``) and the ``filter`` of ``timeseries.search``. + Criteria only: ``limit``, ``sort_by``, ``sort_order`` and ``cursor`` are arguments of the + call, not fields here, so one filter can be reused across ``filter()`` and ``search()`` and + paged differently each time. + ``external_id``, ``name``, ``source``, ``unit`` and ``unit_external_id`` are pattern lists — see ``PatternList``. Each is singular because each also takes a bare string, though a list is always accepted. ``labels`` keeps its plural: its entries must **all** be present, and @@ -239,16 +243,6 @@ class TimeSeriesFilterForm: no restriction, ``[]`` narrows to no datasets and matches nothing. Every other list places no restriction when empty. - ``limit`` defaults to the server's 1000 and is capped at 10000; a value <= 0 falls back to the - default rather than returning nothing. - - ``sort_by`` names one property — ``id``, ``externalId``, ``name``, ``source``, ``description``, - ``createdTime``, ``lastUpdatedTime`` or ``dataSetId`` — with ``sort_order`` of ``"asc"`` or - ``"desc"``; ``id`` is always appended so the order is total. An unrecognised property falls - back to the default, newest created first. Nulls sort last ascending, first descending. - - ``cursor`` continues a previous page: pass that response's ``next_cursor`` verbatim, with the - **same** sort it came from. """ def __init__( self, @@ -264,10 +258,6 @@ class TimeSeriesFilterForm: unit: PatternList | None = None, unit_external_id: PatternList | None = None, value_type: PatternList | None = None, - limit: int | None = None, - sort_by: SortBy | None = None, - sort_order: str | None = None, - cursor: str | None = None, ) -> None: ... @@ -607,7 +597,7 @@ class TimeSeriesServiceSync: def search( self, query: str, - filter: TimeSeriesFilterForm | None = None, + filter: TimeSeriesFilter | None = None, limit: int | None = None, ) -> list[TimeSeries]: """Free-text search for ``query``, ranked by relevance. @@ -617,7 +607,32 @@ class TimeSeriesServiceSync: survives, defaulting to 100 and capping at 1000; the ``filter`` endpoints use 1000/10000, which is easy to conflate. """ - def filter(self, input: TimeSeriesFilterForm) -> Page: ... + def filter( + self, + *, + filter: TimeSeriesFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + unit: PatternList | None = None, + unit_external_id: PatternList | None = None, + value_type: PatternList | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + def insert_datapoints(self, input: list[DatapointsCollectionString]) -> list[str]: ... def insert_from_lists( self, @@ -641,10 +656,35 @@ class TimeSeriesServiceAsync: async def search( self, query: str, - filter: TimeSeriesFilterForm | None = None, + filter: TimeSeriesFilter | None = None, limit: int | None = None, ) -> list[TimeSeries]: ... - async def filter(self, input: TimeSeriesFilterForm) -> Page: ... + async def filter( + self, + *, + filter: TimeSeriesFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + unit: PatternList | None = None, + unit_external_id: PatternList | None = None, + value_type: PatternList | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + async def insert_datapoints(self, input: list[DatapointsCollectionString]) -> list[str]: ... async def insert_from_lists( self, @@ -772,34 +812,6 @@ class EventFilter: ) -> None: ... -class EventFilterForm: - """The ``events.filter`` request body. Omitting ``filter`` places no restriction.""" - def __init__( - self, - filter: EventFilter | None = None, - limit: int | None = None, - sort_by: SortBy | None = None, - sort_order: str | None = None, - cursor: str | None = None, - ) -> None: - """``sort_by`` names one property — ``eventTime``, ``createdTime``, ``lastUpdatedTime``, - ``externalId``, ``type``, ``subType``, ``status``, ``source`` or ``dataSetId``. The default - is ``eventTime`` **ascending**, unlike the node filters' newest-created-first: it is the - order the cursor pages in, so paging does not change the order. - - ``subType`` and ``status`` may be sorted by but **not paged** — they are nullable, and a - keyset boundary on them would skip the events that have no value. - """ - @property - def filter(self) -> EventFilter | None: ... - @property - def limit(self) -> int: ... - @limit.setter - def limit(self, value: int) -> None: ... - @property - def cursor(self) -> str | None: ... - - class EventIdCollection: def __init__( self, @@ -862,7 +874,31 @@ class EventsServiceSync: def get(self, id: UUID) -> Event | None: ... def delete(self, input: list[EventIdentifiable]) -> None: ... def update(self, input: list[EventUpdate]) -> list[Event]: ... - def filter(self, input: EventFilterForm) -> Page: ... + def filter( + self, + *, + filter: EventFilter | None = None, + external_id: PatternList | None = None, + source: PatternList | None = None, + type: PatternList | None = None, + sub_type: PatternList | None = None, + status: PatternList | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + event_time: TimeFilter | None = None, + metadata: MetadataFilter | None = None, + related_resources: Sequence[IdCollection] | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + def search( self, query: str, @@ -892,7 +928,31 @@ class EventsServiceAsync: async def get(self, id: UUID) -> Event | None: ... async def delete(self, input: list[EventIdentifiable]) -> None: ... async def update(self, input: list[EventUpdate]) -> list[Event]: ... - async def filter(self, input: EventFilterForm) -> Page: ... + async def filter( + self, + *, + filter: EventFilter | None = None, + external_id: PatternList | None = None, + source: PatternList | None = None, + type: PatternList | None = None, + sub_type: PatternList | None = None, + status: PatternList | None = None, + data_set_id: Sequence[DataSetRef] | None = None, + event_time: TimeFilter | None = None, + metadata: MetadataFilter | None = None, + related_resources: Sequence[IdCollection] | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + async def search( self, query: str, @@ -1007,18 +1067,6 @@ class DatasetFilter: # `limit` defaults to the server's 1000 and may not exceed 10000. There is no paging, so a filter # broad enough to exceed the cap is truncated — narrow it instead. -class DatasetFilterForm: - def __init__( - self, - filter: DatasetFilter | None = None, - limit: int | None = None, - sort_by: SortBy | None = None, - sort_order: str | None = None, - cursor: str | None = None, - ) -> None: - """See ``TimeSeriesFilterForm`` for the ``sort_by`` / ``sort_order`` / ``cursor`` rules.""" - - # A partial update for one dataset. `dataset` names the target; only the fields you pass are sent, # anything omitted is left untouched. There is deliberately no `policies` or `connected_data_sets` # — the update endpoint does not accept them, whatever a Dataset can carry on create. @@ -1048,7 +1096,28 @@ class DatasetsServiceSync: def create(self, input: list[Dataset]) -> list[Dataset]: ... def by_ids(self, input: list[Identifiable]) -> list[Dataset]: ... def delete(self, input: list[Identifiable]) -> None: ... - def filter(self, input: DatasetFilterForm) -> Page: ... + def filter( + self, + *, + filter: DatasetFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + def search( self, query: str, @@ -1071,7 +1140,28 @@ class DatasetsServiceAsync: async def create(self, input: list[Dataset]) -> list[Dataset]: ... async def by_ids(self, input: list[Identifiable]) -> list[Dataset]: ... async def delete(self, input: list[Identifiable]) -> None: ... - async def filter(self, input: DatasetFilterForm) -> Page: ... + async def filter( + self, + *, + filter: DatasetFilter | None = None, + id: Sequence[int] | None = None, + external_id: PatternList | None = None, + name: PatternList | None = None, + source: PatternList | None = None, + labels: PatternList | None = None, + metadata: MetadataFilter | None = None, + created_time: TimeFilter | None = None, + last_updated_time: TimeFilter | None = None, + limit: int | None = None, + sort_by: SortBy | None = None, + sort_order: str | None = None, + cursor: str | None = None, + ) -> Page: + """Pass either ``filter=`` or the individual criteria keywords; passing both is a + ``TypeError``. Paging is always given here rather than on the filter, so one filter can be + reused across calls. + """ + async def search( self, query: str, @@ -1337,6 +1427,7 @@ class ResourcesServiceSync: def get_by_id(self, id: int) -> Resource | None: ... def filter( self, + filter: ResourceFilter | None = None, id: Sequence[int] | None = None, external_id: PatternList | None = None, name: PatternList | None = None, @@ -1362,7 +1453,7 @@ class ResourcesServiceSync: ``external_id``, ``name`` and ``source`` are pattern lists; ``labels`` must all be present; a ``None`` ``metadata`` value matches the key alone. ``data_set_id`` expands down the dataset hierarchy, and ``None`` (no restriction) differs from ``[]`` - (narrow to no datasets, matching nothing). See ``TimeSeriesFilterForm`` for the sort and + (narrow to no datasets, matching nothing). See ``TimeSeriesFilter`` for the sort and cursor rules. """ @@ -1383,6 +1474,7 @@ class ResourcesServiceAsync: async def get_by_id(self, id: int) -> Resource | None: ... async def filter( self, + filter: ResourceFilter | None = None, id: Sequence[int] | None = None, external_id: PatternList | None = None, name: PatternList | None = None, diff --git a/datahub_python_bindings/src/datasets/async_service.rs b/datahub_python_bindings/src/datasets/async_service.rs index d8a56d4..63efd6f 100644 --- a/datahub_python_bindings/src/datasets/async_service.rs +++ b/datahub_python_bindings/src/datasets/async_service.rs @@ -1,5 +1,5 @@ use crate::datasets::{ - DatasetIdentifiable, PyDataset, PyDatasetFilterForm, PyDatasetUpdate, + DatasetIdentifiable, PyDataset, PyDatasetUpdate, }; use crate::resources::PyResource; use crate::{DatahubIdentity, Identifiable, PyIdCollection}; @@ -109,12 +109,36 @@ impl PyDatasetsServiceAsync { } /// Datasets matching every criterion on the filter, newest first. - fn filter<'p>(&self, py: Python<'p>, input: PyDatasetFilterForm) -> PyResult> { + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=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<'p>( + &self, + py: Python<'p>, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult> { + let form = crate::datasets::dataset_filter_form( + filter, id, external_id, name, source, labels, metadata, 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 .datasets - .filter(&input.into()) + .filter(&form) .await .map_err(crate::datahub_err)?; let next_cursor = result.next_cursor().map(str::to_string); diff --git a/datahub_python_bindings/src/datasets/mod.rs b/datahub_python_bindings/src/datasets/mod.rs index d3f44b1..2590565 100644 --- a/datahub_python_bindings/src/datasets/mod.rs +++ b/datahub_python_bindings/src/datasets/mod.rs @@ -406,7 +406,7 @@ impl PyDatasetFilter { last_updated_time = None, ))] #[allow(clippy::too_many_arguments)] - fn new( + pub fn new( id: Option>, external_id: Option, name: Option, @@ -447,61 +447,59 @@ impl PyDatasetFilter { } } -/// Body of `datasets.filter`: criteria plus a cap. +/// Build the request body for `datasets.filter` from either form of its arguments. +/// +/// 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 -/// rejected. There is no paging, so a filter broad enough to exceed the cap is truncated; -/// narrow it rather than trying to page. -#[pyclass(module = "intellistream_datahub_sdk", name = "DatasetFilterForm", from_py_object)] -#[derive(Clone)] -pub struct PyDatasetFilterForm { - pub inner: DatasetFilterForm, -} - -impl From for PyDatasetFilterForm { - fn from(f: DatasetFilterForm) -> Self { - Self { inner: f } - } -} -impl From for DatasetFilterForm { - fn from(f: PyDatasetFilterForm) -> Self { - f.inner - } -} +/// rejected. +#[allow(clippy::too_many_arguments)] +pub fn dataset_filter_form( + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: 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() + || source.is_some() + || labels.is_some() + || metadata.is_some() + || created_time.is_some() + || last_updated_time.is_some(); + let from_keywords = PyDatasetFilter::new( + id, + external_id, + name, + source, + labels, + metadata, + created_time, + last_updated_time, + ) + .inner; -#[pymethods] -impl PyDatasetFilterForm { - /// The criteria, how many to return, and in what order. - /// - /// `sort_by` names one property — `id`, `externalId`, `name`, `source`, `description`, - /// `createdTime`, `lastUpdatedTime` or `dataSetId` — with `sort_order` of `"asc"` or `"desc"`; - /// `id` is always appended so the order is total. An unrecognised property falls back to the - /// default (newest created first). Nulls sort last ascending, first descending. - /// - /// `cursor` continues a previous page: pass that response's `next_cursor` verbatim, with the - /// **same** sort it came from — a mismatch is a 400, not a quietly short page. - #[new] - #[pyo3(signature = (filter = None, limit = None, sort_by = None, sort_order = None, - cursor = None))] - fn new( - filter: Option, - limit: Option, - sort_by: Option, - sort_order: Option, - cursor: Option, - ) -> Self { - let mut form = DatasetFilterForm::new(); - if let Some(filter) = filter { - form.set_filter(filter.into()); - } - if let Some(limit) = limit { - form.set_limit(limit); - } - form.set_paging(crate::build_page_request(sort_by, sort_order, cursor)); - Self { - inner: form.build(), - } - } + let mut form = DatasetFilterForm::new(); + form.set_filter(crate::resolve_filter( + filter.map(Into::into), + from_keywords, + any_keyword, + )?); + if let Some(limit) = limit { + form.set_limit(limit); + } + form.set_paging(crate::build_page_request(sort_by, sort_order, cursor)); + Ok(form.build()) } /// A partial update for one dataset, mirroring the server's update form. diff --git a/datahub_python_bindings/src/datasets/sync_service.rs b/datahub_python_bindings/src/datasets/sync_service.rs index d8c433d..b4a02bf 100644 --- a/datahub_python_bindings/src/datasets/sync_service.rs +++ b/datahub_python_bindings/src/datasets/sync_service.rs @@ -1,5 +1,5 @@ use crate::datasets::{ - DatasetIdentifiable, PyDataset, PyDatasetFilterForm, PyDatasetUpdate, + DatasetIdentifiable, PyDataset, PyDatasetUpdate, }; use crate::resources::PyResource; use crate::{PyIdCollection}; @@ -90,12 +90,36 @@ impl PyDatasetsServiceSync { } /// Datasets matching every criterion on the filter, newest first. - fn filter(&self, py: Python<'_>, input: PyDatasetFilterForm) -> PyResult { + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=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<'_>, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult { + let form = crate::datasets::dataset_filter_form( + filter, id, external_id, name, source, labels, metadata, created_time, + last_updated_time, limit, sort_by, sort_order, cursor, + )?; let service = self.api_service.clone(); let (items, next_cursor) = py.detach(|| { let result = self .runtime - .block_on(service.datasets.filter(&input.into())) + .block_on(service.datasets.filter(&form)) .map_err(crate::datahub_err)?; let next_cursor = result.next_cursor().map(str::to_string); let items: Vec = result diff --git a/datahub_python_bindings/src/events/async_service.rs b/datahub_python_bindings/src/events/async_service.rs index 897ab7b..d872346 100644 --- a/datahub_python_bindings/src/events/async_service.rs +++ b/datahub_python_bindings/src/events/async_service.rs @@ -1,5 +1,5 @@ use crate::events::{ - EventIdentifyable, PyEventFilter, PyEvent, PyEventDimension, PyEventFilterForm, PyEventUpdate, + EventIdentifyable, PyEventFilter, PyEvent, PyEventDimension, PyEventUpdate, }; use crate::timeseries::async_service::PyTimeSeriesServiceAsync; use crate::timeseries::{PyTimeSeries, PyTimeSeriesUpdate}; @@ -94,13 +94,42 @@ impl PyEventsServiceAsync { }) } - fn filter<'py>(&self, py: Python<'py>, input: PyEventFilterForm) -> PyResult> { + #[pyo3(signature = (filter=None, external_id=None, source=None, r#type=None, sub_type=None, + status=None, data_set_id=None, event_time=None, metadata=None, + related_resources=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>, + filter: Option, + external_id: Option, + source: Option, + r#type: Option, + sub_type: Option, + status: Option, + data_set_id: Option>, + event_time: Option, + metadata: Option>>, + related_resources: Option>, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult> { + let form = crate::events::event_filter_form( + filter, external_id, source, r#type, sub_type, status, data_set_id, event_time, + metadata, related_resources, 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 .events - .filter(&input.into()) + .filter(&form) .await .map_err(|e| crate::datahub_err(e))?; diff --git a/datahub_python_bindings/src/events/mod.rs b/datahub_python_bindings/src/events/mod.rs index 865a1fa..1257433 100644 --- a/datahub_python_bindings/src/events/mod.rs +++ b/datahub_python_bindings/src/events/mod.rs @@ -65,123 +65,73 @@ impl PyEvent { } } -#[pyclass(module = "intellistream_datahub_sdk", name = "EventFilterForm", from_py_object)] -#[derive(Clone)] -pub struct PyEventFilterForm { - pub inner: EventFilterForm, -} -impl From for PyEventFilterForm { - fn from(ts: EventFilterForm) -> Self { - Self { inner: ts } - } -} -impl From for EventFilterForm { - fn from(ts: PyEventFilterForm) -> Self { - ts.inner - } -} - -#[pymethods] -impl PyEventFilterForm { - /// The request body: the criteria, plus paging and ordering. - /// - /// `filter` may be omitted, which places no restriction and returns whatever the tenant - /// has — the same thing an argument-free `EventFilter()` does. - #[new] - #[pyo3(signature=(filter=None,limit=None,sort_by=None,sort_order=None,cursor=None))] - fn new( - filter: Option, - limit: Option, - sort_by: Option, - sort_order: Option, - cursor: Option, - ) -> Self { - let mut form = EventFilterForm::default(); - form.set_filter(filter.map(Into::into).unwrap_or_default()); - form.set_limit(limit.unwrap_or(100)); - // 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 { - form.set_sort(DataSort { - property: property.into(), - order: sort_order, - }); - } - if let Some(cursor) = cursor { - form.set_cursor(cursor); - } - Self { - inner: form.build(), - } - } - #[getter] - fn filter(&self) -> Option { - self.inner.filter().cloned().map(|f| f.into()) - } - #[getter] - pub fn limit(&self) -> u64 { - self.inner.limit - } - #[setter] - pub fn set_limit(&mut self, limit: u64) { - self.inner.limit = limit; - } - #[getter] - pub fn cursor(&self) -> Option<&str> { - self.inner.cursor() - } - - /// Resume a walk from where the previous page stopped: the `next_cursor` of the previous - /// response, verbatim. - /// - /// Opaque — do not build or parse one. Send it with the same `sort_by`/`sort_order` that - /// produced it, or the request is refused with a 400. - #[setter] - pub fn set_cursor(&mut self, cursor: Option) { - match cursor { - Some(cursor) => { - self.inner.set_cursor(cursor); - } - None => { - self.inner.clear_cursor(); - } - } - } - - /// The properties this page is ordered by, if any. - #[getter] - pub fn sort_by(&self) -> Option> { - self.inner.sort().map(|s| s.property.clone()) - } - - #[setter] - pub fn set_sort_by(&mut self, property: Option) { - let order = self.inner.sort().and_then(|s| s.order.clone()); - match property { - Some(property) => { - self.inner.set_sort(DataSort { property: property.into(), order }); - } - None => { - self.inner.clear_sort(); - } - } - } - - /// `"asc"` or `"desc"`. Anything that is not exactly `desc` sorts ascending server-side. - #[getter] - pub fn sort_order(&self) -> Option { - self.inner.sort().and_then(|s| s.order.clone()) - } +/// Build the request body for `events.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 event_filter_form( + filter: Option, + external_id: Option, + source: Option, + r#type: Option, + sub_type: Option, + status: Option, + data_set_id: Option>, + event_time: Option, + metadata: Option>>, + related_resources: Option>, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, +) -> PyResult { + let any_keyword = external_id.is_some() + || source.is_some() + || r#type.is_some() + || sub_type.is_some() + || status.is_some() + || data_set_id.is_some() + || event_time.is_some() + || metadata.is_some() + || related_resources.is_some() + || created_time.is_some() + || last_updated_time.is_some(); + let from_keywords = PyEventFilter::new( + external_id, + source, + r#type, + sub_type, + status, + data_set_id, + event_time, + metadata, + related_resources, + created_time, + last_updated_time, + ) + .inner; - #[setter] - pub fn set_sort_order(&mut self, order: Option) { - let property = self - .inner - .sort() - .map(|s| s.property.clone()) - .unwrap_or_default(); - self.inner.set_sort(DataSort { property, order }); - } + let mut form = EventFilterForm::default(); + form.set_filter(crate::resolve_filter( + filter.map(Into::into), + from_keywords, + any_keyword, + )?); + form.set_limit(limit.unwrap_or(100)); + // 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 { + form.set_sort(DataSort { + property: property.into(), + order: sort_order, + }); + } + if let Some(cursor) = cursor { + form.set_cursor(cursor); + } + Ok(form.build()) } #[pyclass( @@ -191,7 +141,7 @@ impl PyEventFilterForm { )] #[derive(Clone)] pub struct PyEventFilter { - inner: EventFilter, + pub inner: EventFilter, } impl From for PyEventFilter { fn from(ts: EventFilter) -> Self { @@ -238,7 +188,7 @@ impl PyEventFilter { last_updated_time=None, ))] #[allow(clippy::too_many_arguments)] - fn new( + pub fn new( external_id: Option, source: Option, r#type: Option, @@ -504,7 +454,6 @@ 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/events/sync_service.rs b/datahub_python_bindings/src/events/sync_service.rs index 38fe5e2..c2ab01f 100644 --- a/datahub_python_bindings/src/events/sync_service.rs +++ b/datahub_python_bindings/src/events/sync_service.rs @@ -1,5 +1,5 @@ use crate::events::{ - EventIdentifyable, PyEventFilter, PyEvent, PyEventDimension, PyEventFilterForm, PyEventUpdate, + EventIdentifyable, PyEventFilter, PyEvent, PyEventDimension, PyEventUpdate, }; use crate::{PyIdCollection}; use intellistream_datahub_sdk::events::{EventDimension, EventIdCollection, EventUpdate}; @@ -79,13 +79,42 @@ impl PyEventsServiceSync { }) } - fn filter<'py>(&self, py: Python<'py>, input: PyEventFilterForm) -> PyResult { + #[pyo3(signature = (filter=None, external_id=None, source=None, r#type=None, sub_type=None, + status=None, data_set_id=None, event_time=None, metadata=None, + related_resources=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>, + filter: Option, + external_id: Option, + source: Option, + r#type: Option, + sub_type: Option, + status: Option, + data_set_id: Option>, + event_time: Option, + metadata: Option>>, + related_resources: Option>, + created_time: Option, + last_updated_time: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, + ) -> PyResult { + let form = crate::events::event_filter_form( + filter, external_id, source, r#type, sub_type, status, data_set_id, event_time, + metadata, related_resources, created_time, last_updated_time, limit, sort_by, + sort_order, cursor, + )?; let service = self.api_service.clone(); let (items, next_cursor) = py.detach(|| { let result = self .runtime - .block_on(service.events.filter(&input.into())) + .block_on(service.events.filter(&form)) .map_err(|e| crate::datahub_err(e))?; let next_cursor = result.next_cursor().map(str::to_string); diff --git a/datahub_python_bindings/src/lib.rs b/datahub_python_bindings/src/lib.rs index c44a3b6..571eb85 100644 --- a/datahub_python_bindings/src/lib.rs +++ b/datahub_python_bindings/src/lib.rs @@ -564,14 +564,13 @@ impl PyIdCollection { /// The whole paging loop is therefore: /// /// ```python -/// page = client.timeseries.filter(TimeSeriesFilterForm(limit=100, sort_by="name")) +/// page = client.timeseries.filter(limit=100, sort_by="name") /// while page: /// for ts in page: /// ... /// if page.next_cursor is None: /// break -/// page = client.timeseries.filter(TimeSeriesFilterForm( -/// limit=100, sort_by="name", cursor=page.next_cursor)) +/// page = client.timeseries.filter(limit=100, sort_by="name", cursor=page.next_cursor) /// ``` /// /// `next_cursor` is `None` on the last page. A *full* page may still be the last one — the server @@ -741,33 +740,36 @@ pub(crate) fn search_form( } } -#[pyclass(module = "intellistream_datahub_sdk", name = "TimeSeriesFilterForm")] -#[derive(Clone)] -pub struct PyTimeSeriesFilterForm { - pub inner: TimeSeriesFilterForm, +#[pyclass(module = "intellistream_datahub_sdk", name = "TimeSeriesFilter", from_py_object)] +#[derive(Clone, Default)] +pub struct PyTimeSeriesFilter { + pub inner: TimeSeriesFilter, } -impl From for PyTimeSeriesFilterForm { - fn from(form: TimeSeriesFilterForm) -> Self { - Self { inner: form } +impl From for PyTimeSeriesFilter { + fn from(filter: TimeSeriesFilter) -> Self { + Self { inner: filter } } } -impl From for TimeSeriesFilterForm { - fn from(value: PyTimeSeriesFilterForm) -> Self { +impl From for TimeSeriesFilter { + fn from(value: PyTimeSeriesFilter) -> Self { value.inner } } #[pymethods] -impl PyTimeSeriesFilterForm { +impl PyTimeSeriesFilter { /// AND-combined criteria for `timeseries.filter` and the `filter` of `timeseries.search`. /// + /// Criteria only: `limit`, `sort_by`, `sort_order` and `cursor` are arguments of the call, not + /// fields here, so one filter can be reused across `filter()` and `search()` and paged + /// differently each time without carrying a stale cursor into the next use. + /// /// `external_id`, `name`, `source`, `unit` and `unit_external_id` are **pattern** lists: /// `*` and `%` are wildcards, `_` is literal, matching is case-insensitive, and an entry with /// no wildcard matches exactly. Entries within a list OR together; the fields AND. Each of /// them also accepts a bare string. /// /// `labels` must **all** be present. `metadata` entries must all be present too, and a `None` - /// value matches the key alone — `{"health": None}` finds anything tagged `health`, which is - /// what the retired `metadata_key`-without-`metadata_value` used to mean. + /// value matches the key alone — `{"health": None}` finds anything tagged `health`. /// /// `value_type` is matched exactly (case-insensitively) against the closed catalogue /// `BIGINT`, `FLOAT`, `FLOAT32`, `NUMERIC`, `DECIMAL32`, `TEXT`, `MIXED`. @@ -776,15 +778,6 @@ impl PyTimeSeriesFilterForm { /// dataset hierarchy server-side, so a master dataset matches its children's timeseries too. /// **`None` and `[]` differ here**: `None` places no restriction, `[]` narrows to no datasets /// and matches nothing. Every other list is "no restriction" when empty. - /// - /// `sort_by` names one property — `id`, `externalId`, `name`, `source`, `description`, - /// `createdTime`, `lastUpdatedTime` or `dataSetId` — with `sort_order` of `"asc"` or `"desc"`; - /// `id` is always appended so the order is total. An unrecognised property falls back to the - /// default (newest created first) rather than failing. Nulls sort last ascending, first - /// descending. - /// - /// `cursor` continues a previous page: pass the `next_cursor` of that response verbatim, with - /// the **same** sort it came from — a mismatch is a 400, not a quietly short page. #[new] #[pyo3(signature = ( id=None, @@ -799,10 +792,6 @@ impl PyTimeSeriesFilterForm { unit=None, unit_external_id=None, value_type=None, - limit=None, - sort_by=None, - sort_order=None, - cursor=None, ))] #[allow(clippy::too_many_arguments)] pub fn new( @@ -818,36 +807,104 @@ impl PyTimeSeriesFilterForm { unit: Option, unit_external_id: Option, value_type: Option, - limit: Option, - sort_by: Option, - sort_order: Option, - cursor: Option, ) -> Self { Self { - inner: TimeSeriesFilterForm { - filter: TimeSeriesFilter { - node: NodeFilter { - id, - external_id: opt_patterns(external_id), - name: opt_patterns(name), - source: opt_patterns(source), - labels: opt_patterns(labels), - metadata, - created_time: created_time.map(Into::into), - last_updated_time: last_updated_time.map(Into::into), - }, - data_set_id: opt_data_set_refs(data_set_id), - unit: opt_patterns(unit), - unit_external_id: opt_patterns(unit_external_id), - value_type: opt_patterns(value_type), + inner: TimeSeriesFilter { + node: NodeFilter { + id, + external_id: opt_patterns(external_id), + name: opt_patterns(name), + source: opt_patterns(source), + labels: opt_patterns(labels), + metadata, + created_time: created_time.map(Into::into), + last_updated_time: last_updated_time.map(Into::into), }, - limit, - paging: build_page_request(sort_by, sort_order, cursor), + data_set_id: opt_data_set_refs(data_set_id), + unit: opt_patterns(unit), + unit_external_id: opt_patterns(unit_external_id), + value_type: opt_patterns(value_type), }, } } } +/// Build the request body for `timeseries.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 timeseries_filter_form( + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + unit: Option, + unit_external_id: Option, + value_type: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, +) -> PyResult { + let any_keyword = id.is_some() + || external_id.is_some() + || name.is_some() + || source.is_some() + || labels.is_some() + || metadata.is_some() + || created_time.is_some() + || last_updated_time.is_some() + || data_set_id.is_some() + || unit.is_some() + || unit_external_id.is_some() + || value_type.is_some(); + let from_keywords = PyTimeSeriesFilter::new( + id, + external_id, + name, + source, + labels, + metadata, + created_time, + last_updated_time, + data_set_id, + unit, + unit_external_id, + value_type, + ) + .inner; + Ok(TimeSeriesFilterForm { + filter: resolve_filter(filter.map(|f| f.inner), from_keywords, any_keyword)?, + limit, + paging: build_page_request(sort_by, sort_order, cursor), + }) +} + +/// Resolve `filter=` against the flat criteria keywords, for the `filter()` of every service. +/// +/// Passing both is a `TypeError` rather than a merge: a caller who hands over a prepared filter +/// *and* a keyword has two intents in one call, and either answer — keyword wins, or union — +/// silently discards one of them. +pub fn resolve_filter( + filter: Option, + from_keywords: F, + any_keyword_given: bool, +) -> PyResult { + match (filter, any_keyword_given) { + (Some(_), true) => Err(pyo3::exceptions::PyTypeError::new_err( + "pass either filter= or the individual criteria keywords, not both", + )), + (Some(f), false) => Ok(f), + (None, _) => Ok(from_keywords), + } +} + #[derive(FromPyObject)] pub enum Identifiable { #[pyo3(transparent)] @@ -1192,9 +1249,8 @@ fn _core(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::()?; m.add_class::()?; m.add_class::()?; timeseries::register(m)?; diff --git a/datahub_python_bindings/src/resources/async_service.rs b/datahub_python_bindings/src/resources/async_service.rs index 9e6d008..cc8ed9e 100644 --- a/datahub_python_bindings/src/resources/async_service.rs +++ b/datahub_python_bindings/src/resources/async_service.rs @@ -161,14 +161,15 @@ impl PyResourcesServiceAsync { /// `POST /resources/filter` — structured lookup; every criterion is combined with AND. /// See the sync twin for the pattern, label and data-set-scope rules. - #[pyo3(signature = (id=None, external_id=None, name=None, source=None, labels=None, - metadata=None, created_time=None, last_updated_time=None, node_type=None, - is_root=None, data_set_id=None, limit=None, sort_by=None, sort_order=None, - cursor=None))] + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + node_type=None, is_root=None, data_set_id=None, limit=None, sort_by=None, + sort_order=None, cursor=None))] #[allow(clippy::too_many_arguments)] fn filter<'py>( &self, py: Python<'py>, + filter: Option, id: Option>, external_id: Option, name: Option, @@ -186,10 +187,10 @@ impl PyResourcesServiceAsync { cursor: Option, ) -> PyResult> { let form = crate::resources::sync_service::build_resource_filter_form( - id, external_id, name, source, labels, metadata, created_time, + filter, id, external_id, name, source, labels, metadata, created_time, last_updated_time, node_type, is_root, data_set_id, limit, sort_by, sort_order, cursor, - ); + )?; let service = self.api_service.clone(); future_into_py(py, async move { let result = service diff --git a/datahub_python_bindings/src/resources/sync_service.rs b/datahub_python_bindings/src/resources/sync_service.rs index 2ee264c..b3d7418 100644 --- a/datahub_python_bindings/src/resources/sync_service.rs +++ b/datahub_python_bindings/src/resources/sync_service.rs @@ -152,14 +152,15 @@ impl PyResourcesServiceSync { /// `data_set_id` takes numeric ids, external ids, or `IdCollection`s — it used to take ids /// only — and expands down the dataset hierarchy. **`None` and `[]` differ**: `None` places no /// restriction, `[]` narrows to no datasets and matches nothing. - #[pyo3(signature = (id=None, external_id=None, name=None, source=None, labels=None, - metadata=None, created_time=None, last_updated_time=None, node_type=None, - is_root=None, data_set_id=None, limit=None, sort_by=None, sort_order=None, - cursor=None))] + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + node_type=None, is_root=None, data_set_id=None, limit=None, sort_by=None, + sort_order=None, cursor=None))] #[allow(clippy::too_many_arguments)] fn filter<'py>( &self, py: Python<'py>, + filter: Option, id: Option>, external_id: Option, name: Option, @@ -177,10 +178,10 @@ impl PyResourcesServiceSync { cursor: Option, ) -> PyResult { let form = build_resource_filter_form( - id, external_id, name, source, labels, metadata, created_time, + filter, id, external_id, name, source, labels, metadata, created_time, last_updated_time, node_type, is_root, data_set_id, limit, sort_by, sort_order, cursor, - ); + )?; let service = self.api_service.clone(); let (items, next_cursor) = py.detach(|| { let result = self @@ -291,6 +292,7 @@ pub(crate) fn build_resource_filter( /// Shared by the sync and async `filter` bindings: turn Python kwargs into a `ResourceFilterForm`. #[allow(clippy::too_many_arguments)] pub(crate) fn build_resource_filter_form( + filter: Option, id: Option>, external_id: Option, name: Option, @@ -306,16 +308,28 @@ pub(crate) fn build_resource_filter_form( sort_by: Option, sort_order: Option, cursor: Option, -) -> ResourceFilterForm { - let filter = build_resource_filter( +) -> PyResult { + let any_keyword = id.is_some() + || external_id.is_some() + || name.is_some() + || source.is_some() + || labels.is_some() + || metadata.is_some() + || created_time.is_some() + || last_updated_time.is_some() + || node_type.is_some() + || is_root.is_some() + || data_set_id.is_some(); + let from_keywords = build_resource_filter( id, external_id, name, source, labels, metadata, created_time, last_updated_time, node_type, is_root, data_set_id, ); + let filter = crate::resolve_filter(filter.map(|f| f.inner), from_keywords, any_keyword)?; let mut form = ResourceFilterForm::new(filter); if let Some(limit) = limit { form = form.with_limit(limit); } - form.with_paging(crate::build_page_request(sort_by, sort_order, cursor)) + Ok(form.with_paging(crate::build_page_request(sort_by, sort_order, cursor))) } /// Shared by the sync and async `fetch_nearest` bindings. diff --git a/datahub_python_bindings/src/timeseries/async_service.rs b/datahub_python_bindings/src/timeseries/async_service.rs index 29c4fb5..dd794e2 100644 --- a/datahub_python_bindings/src/timeseries/async_service.rs +++ b/datahub_python_bindings/src/timeseries/async_service.rs @@ -6,7 +6,6 @@ use crate::timeseries::{ }; use crate::{ DatahubIdentity, Identifiable, PyIdCollection, PyRetrieveFilter, - PyTimeSeriesFilterForm, }; use crate::datetime::py_datetime_to_utc; use intellistream_datahub_sdk::generic::{ @@ -147,10 +146,10 @@ impl PyTimeSeriesServiceAsync { &self, py: Python<'p>, query: String, - filter: Option, + filter: Option, limit: Option, ) -> PyResult> { - let form = crate::search_form(query, filter.map(|f| f.inner.filter), limit); + let form = crate::search_form(query, filter.map(|f| f.inner), limit); let service = self.api_service.clone(); future_into_py(py, async move { @@ -169,17 +168,43 @@ impl PyTimeSeriesServiceAsync { }) } + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + data_set_id=None, unit=None, unit_external_id=None, value_type=None, + limit=None, sort_by=None, sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] fn filter<'p>( &self, py: Python<'p>, - input: PyTimeSeriesFilterForm, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + unit: Option, + unit_external_id: Option, + value_type: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, ) -> PyResult> { + let form = crate::timeseries_filter_form( + filter, id, external_id, name, source, labels, metadata, created_time, + last_updated_time, data_set_id, unit, unit_external_id, value_type, limit, sort_by, + sort_order, cursor, + )?; let service = self.api_service.clone(); future_into_py(py, async move { let result = service .time_series - .filter(&input.into()) + .filter(&form) .await .map_err(|e| crate::datahub_err(e))?; diff --git a/datahub_python_bindings/src/timeseries/sync_service.rs b/datahub_python_bindings/src/timeseries/sync_service.rs index a3b1184..ef20de3 100644 --- a/datahub_python_bindings/src/timeseries/sync_service.rs +++ b/datahub_python_bindings/src/timeseries/sync_service.rs @@ -4,7 +4,7 @@ use crate::timeseries::datapoints::{ PyDatapointsCollectionDatapoints, PyDatapointsCollectionString, }; use crate::{DatahubIdentity, Identifiable}; -use crate::{PyIdCollection, PyRetrieveFilter, PyTimeSeriesFilterForm}; +use crate::{PyIdCollection, PyRetrieveFilter}; use intellistream_datahub_sdk::generic::{DataWrapper, IdAndExtId}; use intellistream_datahub_sdk::{ApiService, TimeSeriesUpdateCollection}; use pyo3_async_runtimes::tokio::future_into_py; @@ -132,10 +132,10 @@ impl PyTimeSeriesServiceSync { &self, py: Python<'p>, query: String, - filter: Option, + filter: Option, limit: Option, ) -> PyResult> { - let form = crate::search_form(query, filter.map(|f| f.inner.filter), limit); + let form = crate::search_form(query, filter.map(|f| f.inner), limit); let service = self.api_service.clone(); py.detach(|| { @@ -152,17 +152,43 @@ impl PyTimeSeriesServiceSync { }) } + #[pyo3(signature = (filter=None, id=None, external_id=None, name=None, source=None, + labels=None, metadata=None, created_time=None, last_updated_time=None, + data_set_id=None, unit=None, unit_external_id=None, value_type=None, + limit=None, sort_by=None, sort_order=None, cursor=None))] + #[allow(clippy::too_many_arguments)] fn filter<'p>( &self, py: Python<'p>, - input: PyTimeSeriesFilterForm, + filter: Option, + id: Option>, + external_id: Option, + name: Option, + source: Option, + labels: Option, + metadata: Option>>, + created_time: Option, + last_updated_time: Option, + data_set_id: Option>, + unit: Option, + unit_external_id: Option, + value_type: Option, + limit: Option, + sort_by: Option, + sort_order: Option, + cursor: Option, ) -> PyResult { + let form = crate::timeseries_filter_form( + filter, id, external_id, name, source, labels, metadata, created_time, + last_updated_time, data_set_id, unit, unit_external_id, value_type, limit, sort_by, + sort_order, cursor, + )?; let service = self.api_service.clone(); let (items, next_cursor) = py.detach(|| { let result = self .runtime - .block_on(service.time_series.filter(&input.into())) + .block_on(service.time_series.filter(&form)) .map_err(|e| crate::datahub_err(e))?; let next_cursor = result.next_cursor().map(str::to_string); let items: Vec = result diff --git a/python_tests/conftest.py b/python_tests/conftest.py index d99b801..1a7444a 100644 --- a/python_tests/conftest.py +++ b/python_tests/conftest.py @@ -114,11 +114,11 @@ def _sweep(client) -> None: try: events, cursor = [], None while True: - page = client.events.filter(intellistream_datahub_sdk.EventFilterForm( + page = client.events.filter( filter=intellistream_datahub_sdk.EventFilter(external_id=f"{TEST_PREFIX}*"), limit=1000, cursor=cursor, - )) + ) events.extend(page) cursor = getattr(page, "next_cursor", None) if not cursor: diff --git a/python_tests/filter_fixtures.py b/python_tests/filter_fixtures.py index 714bc09..6416212 100644 --- a/python_tests/filter_fixtures.py +++ b/python_tests/filter_fixtures.py @@ -205,8 +205,8 @@ def event_corpus(sync_client, datasets, prefix, token): # Poll rather than sleep: the projection lag is usually milliseconds and occasionally seconds. def visible(): - return sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=50)) + return sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=50) found = poll_until(visible, lambda events: len(events) >= len(specs)) assert len(found) >= len(specs), ( @@ -320,8 +320,8 @@ def sortable_events(sync_client, datasets, prefix, token): ]) def visible(): - return sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), limit=50)) + return sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), limit=50) found = poll_until(visible, lambda events: len(events) >= len(specs)) assert len(found) >= len(specs), "the sortable event corpus never became visible" diff --git a/python_tests/test_datasets.py b/python_tests/test_datasets.py index 3b43d39..f204a9f 100644 --- a/python_tests/test_datasets.py +++ b/python_tests/test_datasets.py @@ -121,26 +121,22 @@ def test_sync_list_and_filter(sync_client, make_dataset): assert any(d.external_id == ext_id for d in sync_client.datasets.list()) narrowed = sync_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm( intellistream_datahub_sdk.DatasetFilter(external_id=[ext_id]) ) - ) assert [d.external_id for d in narrowed] == [ext_id] # An unmatchable criterion is an empty result, not an unfiltered one. assert ( sync_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm( intellistream_datahub_sdk.DatasetFilter(external_id=["ds_does_not_exist_xyz"]) ) - ) == [] ) # An argument-free filter places no restriction, same as list(). assert any( d.external_id == ext_id - for d in sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm()) + for d in sync_client.datasets.filter() ) @@ -149,17 +145,13 @@ def test_sync_filter_by_metadata_and_prefix(sync_client, make_dataset): make_dataset(external_id=ext_id, name=ext_id, metadata={"owner": ext_id}) by_metadata = sync_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm( intellistream_datahub_sdk.DatasetFilter(metadata={"owner": ext_id}) ) - ) assert [d.external_id for d in by_metadata] == [ext_id] by_prefix = sync_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm( intellistream_datahub_sdk.DatasetFilter(external_id=f"{ext_id}*") ) - ) assert [d.external_id for d in by_prefix] == [ext_id] @@ -236,10 +228,8 @@ async def test_async_filter_search_update_policies(async_client, make_dataset): ) narrowed = await async_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm( intellistream_datahub_sdk.DatasetFilter(external_id=[ext_id]) ) - ) assert [d.external_id for d in narrowed] == [ext_id] hits = await async_client.datasets.search(token) diff --git a/python_tests/test_events.py b/python_tests/test_events.py index b94e736..c4a5a71 100644 --- a/python_tests/test_events.py +++ b/python_tests/test_events.py @@ -130,9 +130,9 @@ def test_filter_by_external_id_prefix(sync_client, test_events,event_dataset): # `external_id_prefix` is gone; a trailing `*` asks the same question in the field that # also takes exact ids. filter = intellistream_datahub_sdk.EventFilter(external_id=f"{target_string}*") - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) - results = poll_until(lambda: sync_client.events.filter(filt), lambda r: len(r) >= 1) + results = poll_until(lambda: sync_client.events.filter(**filt), lambda r: len(r) >= 1) assert len(results) >= 1 assert all(e.external_id.startswith(target_string) for e in results) @@ -145,10 +145,10 @@ def test_filter_by_type(sync_client, test_events): filter = intellistream_datahub_sdk.EventFilter( type=[target.type], ) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) results = poll_until( - lambda: sync_client.events.filter(filt), + lambda: sync_client.events.filter(**filt), lambda r: target.external_id in {e.external_id for e in r}, ) assert target.external_id in {e.external_id for e in results} @@ -158,10 +158,10 @@ def test_filter_by_sub_type(sync_client, test_events): filter = intellistream_datahub_sdk.EventFilter( sub_type=[target.sub_type], ) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) results = poll_until( - lambda: sync_client.events.filter(filt), + lambda: sync_client.events.filter(**filt), lambda r: target.external_id in {e.external_id for e in r}, ) assert target.external_id in {e.external_id for e in results} @@ -181,11 +181,11 @@ def test_filter_by_sub_type(sync_client, test_events): def test_filter_by_event_time_range(sync_client, test_events,time_filter,expected_idx): # Events are 1 day apart. Filter for the first 3 days. filter = intellistream_datahub_sdk.EventFilter(event_time=time_filter) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) # The time filter isn't dataset-scoped, so we only assert the window returns something # (poll past ingestion lag) rather than pinning exact membership. - results = poll_until(lambda: sync_client.events.filter(filt), lambda r: len(r) >= 1) + results = poll_until(lambda: sync_client.events.filter(**filt), lambda r: len(r) >= 1) assert len(results) >= 1 @pytest.mark.parametrize("target_idx", [7]) @@ -195,10 +195,10 @@ def test_filter_by_metadata(sync_client, test_events,target_idx): target_metadata = target.metadata filter = intellistream_datahub_sdk.EventFilter(metadata=target_metadata) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) results = poll_until( - lambda: sync_client.events.filter(filt), + lambda: sync_client.events.filter(**filt), lambda r: target.external_id in {e.external_id for e in r}, ) assert target.external_id in {e.external_id for e in results} @@ -218,10 +218,10 @@ def test_filter_by_source(sync_client, test_events): filter = intellistream_datahub_sdk.EventFilter( source=[target.source], ) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) results = poll_until( - lambda: sync_client.events.filter(filt), + lambda: sync_client.events.filter(**filt), lambda r: target.external_id in {e.external_id for e in r}, ) assert target.external_id in {e.external_id for e in results} @@ -229,10 +229,10 @@ def test_filter_by_source(sync_client, test_events): def test_filter_with_limit(sync_client, test_events,event_dataset): filter = intellistream_datahub_sdk.EventFilter(external_id=f"{event_dataset.external_id}*") - # Using the EventFilterForm limit field - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter, limit=5) + # Using the limit argument + filt = dict(filter=filter, limit=5) - results = poll_until(lambda: sync_client.events.filter(filt), lambda r: len(r) == 5) + results = poll_until(lambda: sync_client.events.filter(**filt), lambda r: len(r) == 5) assert len(results) == 5 @@ -241,10 +241,10 @@ def test_filter_by_data_set_ids(sync_client, test_events, event_dataset): # array is rejected with HTTP 400. All fixture events live in event_dataset. target = test_events[0] filter = intellistream_datahub_sdk.EventFilter(data_set_id=[event_dataset.id]) - filt = intellistream_datahub_sdk.EventFilterForm(filter=filter) + filt = dict(filter=filter) results = poll_until( - lambda: sync_client.events.filter(filt), + lambda: sync_client.events.filter(**filt), lambda r: target.external_id in {e.external_id for e in r}, ) assert target.external_id in {e.external_id for e in results} @@ -268,17 +268,17 @@ def test_filter_by_related_resources(sync_client, event_dataset): data_set_id=event_dataset.id, related_resources=[intellistream_datahub_sdk.IdCollection(external_id=res_ext)])]) try: - by_id = intellistream_datahub_sdk.EventFilterForm( + by_id = dict( filter=intellistream_datahub_sdk.EventFilter( related_resources=[intellistream_datahub_sdk.IdCollection(id=res.id)])) - r1 = poll_until(lambda: sync_client.events.filter(by_id), + r1 = poll_until(lambda: sync_client.events.filter(**by_id), lambda r: ev_ext in {e.external_id for e in r}) assert ev_ext in {e.external_id for e in r1}, "filter by related resource id did not find the event" - by_ext = intellistream_datahub_sdk.EventFilterForm( + by_ext = dict( filter=intellistream_datahub_sdk.EventFilter( related_resources=[intellistream_datahub_sdk.IdCollection(external_id=res_ext)])) - r2 = poll_until(lambda: sync_client.events.filter(by_ext), + r2 = poll_until(lambda: sync_client.events.filter(**by_ext), lambda r: ev_ext in {e.external_id for e in r}) assert ev_ext in {e.external_id for e in r2}, "filter by related resource external id did not find the event" @@ -297,18 +297,18 @@ def test_filter_by_created_time(sync_client, test_events, event_dataset): now = pd.Timestamp.now(tz="UTC") day = pd.Timedelta(days=1) - after = intellistream_datahub_sdk.EventFilterForm(filter=intellistream_datahub_sdk.EventFilter( + after = dict(filter=intellistream_datahub_sdk.EventFilter( external_id=f"{event_dataset.external_id}*", created_time=intellistream_datahub_sdk.TimeFilter(start=now - day))) - r = poll_until(lambda: sync_client.events.filter(after), + r = poll_until(lambda: sync_client.events.filter(**after), lambda r: target.external_id in {e.external_id for e in r}) assert target.external_id in {e.external_id for e in r}, "created_time (after) filter did not find the event" # Now that we know it's propagated, the complementary window must exclude it. - before = intellistream_datahub_sdk.EventFilterForm(filter=intellistream_datahub_sdk.EventFilter( + before = dict(filter=intellistream_datahub_sdk.EventFilter( external_id=f"{event_dataset.external_id}*", created_time=intellistream_datahub_sdk.TimeFilter(end=now - day))) - r_before = sync_client.events.filter(before) + r_before = sync_client.events.filter(**before) assert target.external_id not in {e.external_id for e in r_before}, ( "created_time (before yesterday) must not return a just-created event" ) @@ -320,17 +320,17 @@ def test_filter_by_last_updated_time(sync_client, test_events, event_dataset): now = pd.Timestamp.now(tz="UTC") day = pd.Timedelta(days=1) - after = intellistream_datahub_sdk.EventFilterForm(filter=intellistream_datahub_sdk.EventFilter( + after = dict(filter=intellistream_datahub_sdk.EventFilter( external_id=f"{event_dataset.external_id}*", last_updated_time=intellistream_datahub_sdk.TimeFilter(start=now - day))) - r = poll_until(lambda: sync_client.events.filter(after), + r = poll_until(lambda: sync_client.events.filter(**after), lambda r: target.external_id in {e.external_id for e in r}) assert target.external_id in {e.external_id for e in r}, "last_updated_time (after) filter did not find the event" - before = intellistream_datahub_sdk.EventFilterForm(filter=intellistream_datahub_sdk.EventFilter( + before = dict(filter=intellistream_datahub_sdk.EventFilter( external_id=f"{event_dataset.external_id}*", last_updated_time=intellistream_datahub_sdk.TimeFilter(end=now - day))) - r_before = sync_client.events.filter(before) + r_before = sync_client.events.filter(**before) assert target.external_id not in {e.external_id for e in r_before}, ( "last_updated_time (before yesterday) must not return a just-created event" ) diff --git a/python_tests/test_filter_datasets.py b/python_tests/test_filter_datasets.py index a2c17e2..91bc97c 100644 --- a/python_tests/test_filter_datasets.py +++ b/python_tests/test_filter_datasets.py @@ -37,8 +37,8 @@ def flt(sync_client, prefix): """Filter within this run's corpus unless the test overrides ``external_id``.""" def _filter(limit=None, **criteria): criteria.setdefault("external_id", f"{prefix}*") - form = intellistream_datahub_sdk.DatasetFilterForm(intellistream_datahub_sdk.DatasetFilter(**criteria), limit=limit) - return externals(sync_client.datasets.filter(form)) + form = dict(filter=intellistream_datahub_sdk.DatasetFilter(**criteria), limit=limit) + return externals(sync_client.datasets.filter(**form)) return _filter @@ -160,7 +160,7 @@ def test_the_retired_flags_are_not_accepted(sync_client): def test_an_argument_free_filter_returns_everything(sync_client, datasets, both): """The same thing ``/datasets/list`` does — the server implements ``list`` by calling the filter handler with an empty filter.""" - from_filter = externals(sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm())) + from_filter = externals(sync_client.datasets.filter()) assert from_filter >= both assert externals(sync_client.datasets.list(limit=10_000)) >= both @@ -188,9 +188,9 @@ def test_an_unmatchable_criterion_gives_an_empty_result_not_an_unfiltered_one(fl # --------------------------------------------------------------------------- # def test_limit_caps_the_page(sync_client, datasets, prefix): - form = intellistream_datahub_sdk.DatasetFilterForm( - intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*"), limit=1) - assert len(sync_client.datasets.filter(form)) == 1 + form = dict( + filter=intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*"), limit=1) + assert len(sync_client.datasets.filter(**form)) == 1 def test_limit_zero_falls_back_to_the_default(flt, both): @@ -198,10 +198,10 @@ def test_limit_zero_falls_back_to_the_default(flt, both): def test_a_limit_above_the_ceiling_is_refused(sync_client, prefix): - form = intellistream_datahub_sdk.DatasetFilterForm( - intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*"), limit=10_001) + form = dict( + filter=intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*"), limit=10_001) with pytest.raises(DataHubException) as excinfo: - sync_client.datasets.filter(form) + sync_client.datasets.filter(**form) assert excinfo.value.status_code == 400 @@ -217,6 +217,6 @@ def test_the_retired_criteria_are_not_accepted(sync_client): @pytest.mark.asyncio async def test_async_filter_matches_the_sync_one(async_client, sync_client, datasets, prefix): - form = intellistream_datahub_sdk.DatasetFilterForm(intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*")) - assert externals(await async_client.datasets.filter(form)) == externals( - sync_client.datasets.filter(form)) + form = dict(filter=intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}*")) + assert externals(await async_client.datasets.filter(**form)) == externals( + sync_client.datasets.filter(**form)) diff --git a/python_tests/test_filter_events.py b/python_tests/test_filter_events.py index 2bdf193..982bec4 100644 --- a/python_tests/test_filter_events.py +++ b/python_tests/test_filter_events.py @@ -50,11 +50,11 @@ def flt(sync_client, prefix): """ def _filter(limit=None, expect_rows=True, **criteria): criteria.setdefault("external_id", f"{prefix}*") - request = intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(**criteria), limit=limit or 100) + request = dict( + filter=intellistream_datahub_sdk.EventFilter(**criteria), limit=limit or 100) def fetch(): - return externals(sync_client.events.filter(request)) + return externals(sync_client.events.filter(**request)) if not expect_rows: return fetch() @@ -231,8 +231,8 @@ def test_empty_and_absent_data_set_scopes_are_opposites(sync_client, event_corpu """``None`` is "no data set restriction", ``[]`` is "narrow to no data sets". Opposite answers, so the SDK has to keep them apart all the way onto the wire.""" def run(data_set_ids): - return externals(sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*", data_set_id=data_set_ids)))) + return externals(sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*", data_set_id=data_set_ids))) assert run([]) == set() assert poll_until(lambda: run(None), bool, timeout=10.0) == both @@ -286,9 +286,9 @@ def test_related_resources_must_all_be_attached(sync_client, datasets, token): ]) try: def run(related): - return externals(sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( + return externals(sync_client.events.filter( intellistream_datahub_sdk.EventFilter(external_id=f"{own_prefix}_ev_*", - related_resources=related)))) + related_resources=related))) by_a = [intellistream_datahub_sdk.IdCollection(external_id=res_a)] assert poll_until(lambda: run(by_a), lambda r: len(r) >= 2, timeout=15.0) == {ev_both, ev_one} @@ -346,8 +346,7 @@ def test_created_time_is_ingest_time_not_event_time(flt, event_corpus, both): def test_an_absent_filter_places_no_restriction(sync_client, event_corpus, both): """An argument-free filter returns the tenant's events, not none of them.""" everything = poll_until( - lambda: externals(sync_client.events.filter( - intellistream_datahub_sdk.EventFilterForm(limit=1000))), + lambda: externals(sync_client.events.filter(limit=1000)), lambda found: found >= both, timeout=15.0, ) @@ -357,9 +356,9 @@ def test_an_absent_filter_places_no_restriction(sync_client, event_corpus, both) @pytest.mark.parametrize("empty", [[], ["", " "]]) def test_empty_and_blank_lists_place_no_restriction(sync_client, event_corpus, prefix, both, empty): scoped = poll_until( - lambda: externals(sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( + lambda: externals(sync_client.events.filter( intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*", type=empty, - sub_type=empty, status=empty, source=empty)))), + sub_type=empty, status=empty, source=empty))), bool, timeout=10.0, ) @@ -377,10 +376,10 @@ def test_garbled_criteria_match_nothing_without_erroring(flt, event_corpus): def test_a_bare_string_means_a_one_element_list(sync_client, event_corpus, prefix, alarm, token): - scalar = sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*", type=f"alarm_{token}"))) - listed = sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=[f"{prefix}*"], type=[f"alarm_{token}"]))) + scalar = sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*", type=f"alarm_{token}")) + listed = sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=[f"{prefix}*"], type=[f"alarm_{token}"])) assert externals(scalar) == externals(listed) == {alarm} @@ -390,8 +389,8 @@ def test_a_bare_string_means_a_one_element_list(sync_client, event_corpus, prefi def test_limit_caps_the_page(sync_client, event_corpus, prefix): capped = poll_until( - lambda: sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=1)), + lambda: sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=1), bool, timeout=10.0, ) @@ -400,8 +399,8 @@ def test_limit_caps_the_page(sync_client, event_corpus, prefix): def test_a_limit_above_the_ceiling_is_refused(sync_client, prefix): with pytest.raises(DataHubException) as excinfo: - sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=10_001)) + sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=10_001) assert excinfo.value.status_code == 400 @@ -409,8 +408,8 @@ def _ordered_page(sync_client, prefix, both, **sort): return [ event.external_id for event in poll_until( - lambda: sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=100, **sort)), + lambda: sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*"), limit=100, **sort), lambda found: externals(found) >= both, timeout=10.0, ) @@ -464,8 +463,8 @@ def test_the_default_order_is_event_time_then_id_ascending(sync_client, datasets try: expected = [externals_by_offset[i] for i in range(6)] page = poll_until( - lambda: sync_client.events.filter(intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{own_prefix}*"), limit=100)), + lambda: sync_client.events.filter( + intellistream_datahub_sdk.EventFilter(external_id=f"{own_prefix}*"), limit=100), lambda found: externals(found) == set(expected), timeout=20.0, ) @@ -492,6 +491,6 @@ def test_the_retired_criteria_are_not_accepted(sync_client): @pytest.mark.asyncio async def test_async_filter_matches_the_sync_one(async_client, sync_client, event_corpus, prefix): - request = intellistream_datahub_sdk.EventFilterForm(intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*")) - assert externals(await async_client.events.filter(request)) == externals( - sync_client.events.filter(request)) + request = dict(filter=intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}*")) + assert externals(await async_client.events.filter(**request)) == externals( + sync_client.events.filter(**request)) diff --git a/python_tests/test_filter_paging.py b/python_tests/test_filter_paging.py index b061f2c..4f49026 100644 --- a/python_tests/test_filter_paging.py +++ b/python_tests/test_filter_paging.py @@ -60,8 +60,7 @@ def forge_cursor(property_name, direction, row_id, value): def ts_page(sync_client, prefix): """One page of the sortable timeseries corpus.""" def _page(**paging): - return sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", **paging)) + return sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", **paging) return _page @@ -108,8 +107,7 @@ def test_a_page_behaves_like_a_list(ts_page, by_index): def test_an_empty_page_is_falsey(sync_client, prefix): - page = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id="no_such_external_id_at_all")) + page = sync_client.timeseries.filter(external_id="no_such_external_id_at_all") assert len(page) == 0 assert bool(page) is False assert page == [] @@ -164,8 +162,7 @@ def test_paging_follows_the_sort_direction(ts_page, by_index): def test_paging_the_default_order(ts_page, sync_client, prefix): """No explicit sort, so the walk runs in the default newest-created-first order.""" rows, _requests = walk(ts_page, limit=2) - expected = ids_of(sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", limit=100))) + expected = ids_of(sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", limit=100)) assert rows == expected @@ -180,8 +177,8 @@ def test_paging_across_a_run_of_tied_values(ts_page, sync_client, prefix): assert len(rows) == 6 assert len(set(rows)) == 6, f"a tied boundary repeated or dropped rows: {rows}" - unpaged = ids_of(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="asc", limit=100))) + unpaged = ids_of(sync_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="asc", limit=100)) assert rows == unpaged @@ -234,13 +231,10 @@ def test_every_filter_endpoint_rejects_an_unreadable_cursor(sync_client, prefix, reported for a bad request. One shape now, from one ``@ControllerAdvice``. """ calls = { - "timeseries": lambda: sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(limit=2, cursor="not-a-cursor")), + "timeseries": lambda: sync_client.timeseries.filter(limit=2, cursor="not-a-cursor"), "resources": lambda: sync_client.resources.filter(limit=2, cursor="not-a-cursor"), - "datasets": lambda: sync_client.datasets.filter( - intellistream_datahub_sdk.DatasetFilterForm(limit=2, cursor="not-a-cursor")), - "events": lambda: sync_client.events.filter( - intellistream_datahub_sdk.EventFilterForm(limit=2, cursor="not-a-cursor")), + "datasets": lambda: sync_client.datasets.filter(limit=2, cursor="not-a-cursor"), + "events": lambda: sync_client.events.filter(limit=2, cursor="not-a-cursor"), } for endpoint, call in calls.items(): with pytest.raises(DataHubException) as excinfo: @@ -286,8 +280,8 @@ def test_a_malformed_boundary_is_a_400_and_not_a_500(sync_client, prefix, sortab for property_name in ["id", "dataSetId", "createdTime", "lastUpdatedTime"]: cursor = forge_cursor(property_name, "asc", "5", "not-a-number") with pytest.raises(DataHubException) as excinfo: - sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", limit=2, sort_by=property_name, cursor=cursor)) + sync_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", limit=2, sort_by=property_name, cursor=cursor) assert excinfo.value.status_code == 400, \ f"{property_name}: {excinfo.value.status_code} {excinfo.value.message[:100]}" @@ -302,17 +296,17 @@ def test_an_injection_payload_in_the_cursor_boundary_is_data(sync_client, prefix distinguishes data from syntax; "it returned no rows" would not. """ def page_after(boundary): - return ids_of(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( + return ids_of(sync_client.timeseries.filter( external_id=f"{prefix}_sort_ts_*", limit=10, sort_by="name", sort_order="asc", - cursor=forge_cursor("name", "asc", "1", boundary)))) + cursor=forge_cursor("name", "asc", "1", boundary))) ordinary = page_after("!") for payload in ["' OR 1=1 --", "'; DROP TABLE node;--", "' UNION SELECT 1 --", "{x:String}"]: assert page_after(payload) == ordinary, f"payload={payload!r}" # And the table is still there afterwards. - assert ids_of(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", sort_by="name", sort_order="asc"))) == by_index + assert ids_of(sync_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", sort_by="name", sort_order="asc")) == by_index # --------------------------------------------------------------------------- # @@ -370,9 +364,9 @@ def test_datasets_page(sync_client, datasets, prefix): parent, child = datasets def page(cursor=None): - return sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm( + return sync_client.datasets.filter( intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}_ds_*"), - limit=1, sort_by="name", sort_order="asc", cursor=cursor)) + limit=1, sort_by="name", sort_order="asc", cursor=cursor) rows, requests = walk(lambda cursor=None, **_: page(cursor)) assert rows == [child.external_id, parent.external_id] @@ -398,9 +392,9 @@ def page(cursor=None): @pytest.fixture def ev_page(sync_client, prefix, sortable_events): def _page(**paging): - request = intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), **paging) - return sync_client.events.filter(request) + request = dict( + filter=intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), **paging) + return sync_client.events.filter(**request) return _page @@ -480,12 +474,12 @@ def test_an_exhausted_event_walk_omits_the_cursor_rather_than_nulling_it(ev_page @pytest.mark.asyncio async def test_async_paging_matches_the_sync_client(async_client, sync_client, prefix, sortable_timeseries, by_index): - first = await async_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", limit=2, sort_by="name", sort_order="asc")) + first = await async_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", limit=2, sort_by="name", sort_order="asc") assert ids_of(first) == by_index[:2] assert first.next_cursor is not None - second = await async_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( + second = await async_client.timeseries.filter( external_id=f"{prefix}_sort_ts_*", limit=2, sort_by="name", sort_order="asc", - cursor=first.next_cursor)) + cursor=first.next_cursor) assert ids_of(second) == by_index[2:4] diff --git a/python_tests/test_filter_search_bodies.py b/python_tests/test_filter_search_bodies.py index 36b90fe..e5f4ba2 100644 --- a/python_tests/test_filter_search_bodies.py +++ b/python_tests/test_filter_search_bodies.py @@ -61,7 +61,7 @@ def test_timeseries_search_filter_narrows_on_an_inherited_node_field(sync_client narrowed = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(name=[f"Pump Alpha {token}"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(name=[f"Pump Alpha {token}"]), ) assert externals(narrowed) == {timeseries_corpus["pump_1"].external_id} @@ -70,13 +70,13 @@ def test_timeseries_search_filter_narrows_on_a_timeseries_only_field(sync_client query = f"Pump {token}" by_unit = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(unit=["celsius"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(unit=["celsius"]), ) assert externals(by_unit) == {timeseries_corpus["pump_x1"].external_id} by_value_type = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(value_type=["TEXT"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(value_type=["TEXT"]), ) assert externals(by_value_type) == set(), "both Pump series are FLOAT" @@ -85,13 +85,13 @@ def test_timeseries_search_filter_narrows_by_metadata_and_labels(sync_client, ti query = f"Pump {token}" by_metadata = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(metadata={f"tsk_{token}": "beta"}), + filter=intellistream_datahub_sdk.TimeSeriesFilter(metadata={f"tsk_{token}": "beta"}), ) assert externals(by_metadata) == {timeseries_corpus["pump_x1"].external_id} by_label = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(labels=["NO_SUCH_LABEL_XYZ"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(labels=["NO_SUCH_LABEL_XYZ"]), ) assert externals(by_label) == set() @@ -107,14 +107,14 @@ def test_timeseries_search_filter_narrows_by_data_set(sync_client, timeseries_co in_child = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(data_set_id=[child.id]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(data_set_id=[child.id]), ) assert timeseries_corpus["valve"].external_id not in externals(in_child) # Naming the parent covers the child, so the whole corpus is back in scope. under_parent = sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(data_set_id=[parent.id]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(data_set_id=[parent.id]), ) assert externals(under_parent) >= { timeseries_corpus["pump_1"].external_id, timeseries_corpus["valve"].external_id @@ -141,7 +141,7 @@ def test_timeseries_search_ranking_survives_the_filter(sync_client, timeseries_c )] filtered = [ts.external_id for ts in sync_client.timeseries.search( query, - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(value_type=["FLOAT"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(value_type=["FLOAT"]), )] kept = [external_id for external_id in unfiltered if external_id in set(filtered)] @@ -154,7 +154,7 @@ async def test_timeseries_search_filter_works_on_the_async_client( ): narrowed = await async_client.timeseries.search( f"Pump {token}", - filter=intellistream_datahub_sdk.TimeSeriesFilterForm(name=[f"Pump Alpha {token}"]), + filter=intellistream_datahub_sdk.TimeSeriesFilter(name=[f"Pump Alpha {token}"]), ) assert externals(narrowed) == {timeseries_corpus["pump_1"].external_id} diff --git a/python_tests/test_filter_sorting.py b/python_tests/test_filter_sorting.py index 9cb2866..7d2adfc 100644 --- a/python_tests/test_filter_sorting.py +++ b/python_tests/test_filter_sorting.py @@ -45,8 +45,7 @@ def ids_of(page): def ts_sorted(sync_client, prefix): """The sortable timeseries corpus, in the order the server returns it.""" def _sorted(**paging): - return ids_of(sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", **paging))) + return ids_of(sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", **paging)) return _sorted @@ -75,23 +74,20 @@ def test_sorting_by_external_id(ts_sorted, by_index): def test_sorting_by_id(ts_sorted, by_index, sync_client, prefix): """``id`` is both a sortable property and the implicit tie-breaker. Ids are assigned in creation order, which is deliberately not index order — so this is a third distinct sequence.""" - rows = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", sort_by="id", - sort_order="asc")) + rows = sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", sort_by="id", + sort_order="asc") numeric = [ts.id for ts in rows] assert numeric == sorted(numeric) assert set(ids_of(rows)) == set(by_index) - descending = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", sort_by="id", - sort_order="desc")) + descending = sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", sort_by="id", + sort_order="desc") assert [ts.id for ts in descending] == sorted(numeric, reverse=True) def test_sorting_by_created_time_is_creation_order(ts_sorted, sync_client, prefix): - rows = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", - sort_by="createdTime", sort_order="asc")) + rows = sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", + sort_by="createdTime", sort_order="asc") stamps = [ts.created_time for ts in rows] assert stamps == sorted(stamps), f"not ascending by createdTime: {stamps}" assert len(stamps) == 6 @@ -99,11 +95,9 @@ def test_sorting_by_created_time_is_creation_order(ts_sorted, sync_client, prefi def test_the_default_order_is_newest_created_first(ts_sorted, sync_client, prefix): """What the node filters returned before they could be sorted, kept as the default.""" - unsorted = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*")) - explicit = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*", - sort_by="createdTime", sort_order="desc")) + unsorted = sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*") + explicit = sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*", + sort_by="createdTime", sort_order="desc") assert ids_of(unsorted) == ids_of(explicit) stamps = [ts.created_time for ts in unsorted] @@ -159,14 +153,13 @@ def test_ties_are_broken_by_id(sync_client, prefix, sortable_timeseries): Sorted by ``dataSetId``, which is identical for all six. """ def page(): - return ids_of(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="asc"))) + return ids_of(sync_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="asc")) first, second = page(), page() assert first == second, "a tied sort must still be deterministic" - by_id = {ts.external_id: ts.id for ts in sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*"))} + by_id = {ts.external_id: ts.id for ts in sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*")} assert [by_id[external_id] for external_id in first] == sorted(by_id.values()), \ "ties should fall back to id ascending" @@ -174,10 +167,9 @@ def page(): def test_the_tie_breaker_follows_the_sort_direction(sync_client, prefix, sortable_timeseries): """Descending means descending all the way down, or the two halves of the order disagree and a keyset boundary lands in the wrong place.""" - descending = ids_of(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="desc"))) - by_id = {ts.external_id: ts.id for ts in sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}_sort_ts_*"))} + descending = ids_of(sync_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", sort_by="dataSetId", sort_order="desc")) + by_id = {ts.external_id: ts.id for ts in sync_client.timeseries.filter(external_id=f"{prefix}_sort_ts_*")} assert [by_id[external_id] for external_id in descending] == sorted( by_id.values(), reverse=True) @@ -218,9 +210,9 @@ def test_datasets_sort_by_name(sync_client, datasets, prefix, token): parent, child = datasets def page(order): - return ids_of(sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm( + return ids_of(sync_client.datasets.filter( intellistream_datahub_sdk.DatasetFilter(external_id=f"{prefix}_ds_*"), - sort_by="name", sort_order=order))) + sort_by="name", sort_order=order)) # "Filter Child" sorts before "Filter Parent". assert page("asc") == [child.external_id, parent.external_id] @@ -245,10 +237,10 @@ def page(order): @pytest.fixture def ev_sorted(sync_client, prefix, sortable_events): def _sorted(**paging): - request = intellistream_datahub_sdk.EventFilterForm( - intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), limit=50, **paging) + request = dict( + filter=intellistream_datahub_sdk.EventFilter(external_id=f"{prefix}_sort_ev_*"), limit=50, **paging) return poll_until( - lambda: ids_of(sync_client.events.filter(request)), + lambda: ids_of(sync_client.events.filter(**request)), lambda found: len(found) >= 4, timeout=15.0, ) @@ -316,6 +308,6 @@ def test_events_sort_by_a_nullable_property(ev_sorted, sortable_events): @pytest.mark.asyncio async def test_async_sorting_matches_the_sync_client(async_client, sync_client, prefix, sortable_timeseries, by_index): - from_async = await async_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}_sort_ts_*", sort_by="name", sort_order="asc")) + from_async = await async_client.timeseries.filter( + external_id=f"{prefix}_sort_ts_*", sort_by="name", sort_order="asc") assert ids_of(from_async) == by_index diff --git a/python_tests/test_filter_timeseries.py b/python_tests/test_filter_timeseries.py index 710a02b..4495745 100644 --- a/python_tests/test_filter_timeseries.py +++ b/python_tests/test_filter_timeseries.py @@ -47,7 +47,7 @@ def flt(sync_client, prefix): """ def _filter(**criteria): criteria.setdefault("external_id", f"{prefix}*") - return externals(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm(**criteria))) + return externals(sync_client.timeseries.filter(**criteria)) return _filter @@ -131,10 +131,8 @@ def test_a_bare_string_means_a_one_element_list(sync_client, timeseries_corpus, Filter fields went plural because one value was rarely enough, but most calls still pass one; making the single form a `TypeError` would tax the common case to serve the rare one. """ - scalar = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", name=f"Pump Alpha {token}")) - listed = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=[f"{prefix}*"], name=[f"Pump Alpha {token}"])) + scalar = sync_client.timeseries.filter(external_id=f"{prefix}*", name=f"Pump Alpha {token}") + listed = sync_client.timeseries.filter(external_id=[f"{prefix}*"], name=[f"Pump Alpha {token}"]) assert externals(scalar) == externals(listed) == {timeseries_corpus["pump_1"].external_id} @@ -244,17 +242,17 @@ def test_an_unknown_value_type_matches_nothing(flt, timeseries_corpus): def test_an_absent_filter_places_no_restriction(sync_client, timeseries_corpus): """An argument-free form must return the tenant's timeseries, not none of them.""" - everything = sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm()) + everything = sync_client.timeseries.filter() assert externals(everything) >= {ts.external_id for ts in timeseries_corpus.values()} def test_none_valued_criteria_are_omitted_entirely(sync_client, timeseries_corpus, prefix): """Passing ``None`` is the same as not passing the argument — it must not reach the wire as a ``null`` the server then reads as a restriction.""" - assert externals(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( + assert externals(sync_client.timeseries.filter( external_id=f"{prefix}*", name=None, unit=None, value_type=None, metadata=None, id=None, labels=None, source=None, data_set_id=None, - ))) == {ts.external_id for ts in timeseries_corpus.values()} + )) == {ts.external_id for ts in timeseries_corpus.values()} @pytest.mark.parametrize("empty", [[], ["", " "]]) @@ -265,8 +263,8 @@ def test_an_empty_or_blank_list_places_no_restriction(sync_client, timeseries_co ``data_set_id`` is the documented exception and is covered separately. """ - scoped = externals(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}*", name=empty, unit=empty, labels=empty, value_type=empty))) + scoped = externals(sync_client.timeseries.filter( + external_id=f"{prefix}*", name=empty, unit=empty, labels=empty, value_type=empty)) assert scoped == {ts.external_id for ts in timeseries_corpus.values()} @@ -309,8 +307,8 @@ def test_an_explicit_empty_data_set_scope_matches_nothing(sync_client, timeserie dropping the predicate instead would widen the query to everything the caller can read — the opposite of what they asked for. """ - assert externals(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}*", data_set_id=[]))) == set() + assert externals(sync_client.timeseries.filter( + external_id=f"{prefix}*", data_set_id=[])) == set() def test_an_omitted_data_set_scope_places_no_restriction(flt, timeseries_corpus): @@ -383,23 +381,21 @@ def test_separate_criteria_and_together(flt, timeseries_corpus, token): def test_limit_caps_the_page(sync_client, timeseries_corpus, prefix): - capped = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", limit=2)) + capped = sync_client.timeseries.filter(external_id=f"{prefix}*", limit=2) assert len(capped) == 2 def test_limit_zero_falls_back_to_the_default(sync_client, timeseries_corpus, prefix): """SQL reads ``LIMIT 0`` as "return nothing", which is indistinguishable from "nothing matched" — so the server treats a non-positive limit as unset instead.""" - assert externals(sync_client.timeseries.filter(intellistream_datahub_sdk.TimeSeriesFilterForm( - external_id=f"{prefix}*", limit=0))) == {ts.external_id for ts in timeseries_corpus.values()} + assert externals(sync_client.timeseries.filter( + external_id=f"{prefix}*", limit=0)) == {ts.external_id for ts in timeseries_corpus.values()} def test_a_limit_above_the_ceiling_is_refused(sync_client, prefix): """10000 is the cap, and exceeding it is a 400 rather than a silently clamped page.""" with pytest.raises(DataHubException) as excinfo: - sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", limit=10_001)) + sync_client.timeseries.filter(external_id=f"{prefix}*", limit=10_001) assert excinfo.value.status_code == 400 @@ -409,14 +405,11 @@ def test_a_negative_limit_is_rejected_client_side(sync_client, prefix): difference between "the server tolerates it" and "you cannot send it" stays visible. """ with pytest.raises(OverflowError): - sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", limit=-5)) + sync_client.timeseries.filter(external_id=f"{prefix}*", limit=-5) @pytest.mark.asyncio async def test_async_filter_matches_the_sync_one(async_client, sync_client, timeseries_corpus, prefix): - from_async = await async_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", unit=["bar"])) - from_sync = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(external_id=f"{prefix}*", unit=["bar"])) + from_async = await async_client.timeseries.filter(external_id=f"{prefix}*", unit=["bar"]) + from_sync = sync_client.timeseries.filter(external_id=f"{prefix}*", unit=["bar"]) assert externals(from_async) == externals(from_sync) diff --git a/python_tests/test_timeseries_search.py b/python_tests/test_timeseries_search.py index a830a35..aa7fce0 100644 --- a/python_tests/test_timeseries_search.py +++ b/python_tests/test_timeseries_search.py @@ -84,14 +84,12 @@ def test_a_name_is_matched_through_the_filter(sync_client, make_ts): time.sleep(SEARCH_INDEX_DELAY) - exact = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(name=unique_name)) + exact = sync_client.timeseries.filter(name=unique_name) assert any(t.external_id == ext_id for t in exact), ( f"filtering by name did not return {ext_id}" ) - as_pattern = sync_client.timeseries.filter( - intellistream_datahub_sdk.TimeSeriesFilterForm(name=f"Py SDK Search {token[:6]}*")) + as_pattern = sync_client.timeseries.filter(name=f"Py SDK Search {token[:6]}*") assert any(t.external_id == ext_id for t in as_pattern), ( "the name filter is a pattern list, so a trailing wildcard must match" ) diff --git a/python_tests/test_update_datasets.py b/python_tests/test_update_datasets.py index 9ab68d9..df5806f 100644 --- a/python_tests/test_update_datasets.py +++ b/python_tests/test_update_datasets.py @@ -78,9 +78,9 @@ def test_external_id_set_value(sync_client, new_dataset): )) found = poll_until( - lambda: sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm( + lambda: sync_client.datasets.filter( intellistream_datahub_sdk.DatasetFilter(external_id=[new_ext]) - )), + ), bool, ) assert [d.external_id for d in found] == [new_ext] @@ -263,9 +263,9 @@ def test_update_persists_beyond_the_echo(sync_client, new_dataset): )) stored = poll_until( - lambda: sync_client.datasets.filter(intellistream_datahub_sdk.DatasetFilterForm( + lambda: sync_client.datasets.filter( intellistream_datahub_sdk.DatasetFilter(external_id=[dataset.external_id]) - )), + ), lambda found: any(d.description == "after" for d in found), ) assert any(d.description == "after" for d in stored)