Skip to content
29 changes: 28 additions & 1 deletion cognite/client/_api/data_modeling/time_series.py
Original file line number Diff line number Diff line change
Expand Up @@ -9,7 +9,8 @@
from cognite.client.data_classes.data_modeling.ids import NodeId, ViewId
from cognite.client.data_classes.data_modeling.instances import InstanceSort, Node, NodeList
from cognite.client.data_classes.data_modeling.views import View
from cognite.client.data_classes.filters import Filter
from cognite.client.data_classes.filters import Equals, Filter, In
from cognite.client.data_classes.time_series import TimeSeriesType
from cognite.client.utils._data_modeling import resolve_source, strip_canonical_source
from cognite.client.utils.useful_types import SequenceNotStr

Expand All @@ -21,6 +22,24 @@
COGNITE_TIME_SERIES_VIEW_ID = CogniteTimeSeries.get_source()


def _build_filter(
filter: Filter | dict[str, Any] | None, time_series_type: TimeSeriesType | Sequence[TimeSeriesType] | None = None
) -> Filter | None:
if isinstance(filter, dict):
filter = Filter.load(filter)

if time_series_type is None:
return filter

type_prop = COGNITE_TIME_SERIES_VIEW_ID.as_property_ref("type")
types = [time_series_type] if isinstance(time_series_type, str) else list(time_series_type)
if not types:
raise ValueError("'time_series_type' must not be empty, pass None to list all kinds of time series")

type_flt = Equals(type_prop, value=types[0]) if len(types) == 1 else In(type_prop, values=types)
return type_flt if filter is None else type_flt & filter


class DataModelingTimeSeriesAPI(APIClient):
def __init__(self, config: ClientConfig, api_version: str | None, cognite_client: AsyncCogniteClient) -> None:
# TODO: Add DataModelingDatapointsAPI
Expand Down Expand Up @@ -107,6 +126,7 @@ async def list(
space: str | SequenceNotStr[str] | None = None,
sort: Sequence[InstanceSort | dict] | InstanceSort | dict | None = None,
filter: Filter | dict[str, Any] | None = None,
time_series_type: TimeSeriesType | Sequence[TimeSeriesType] | None = None,
limit: int | None = DEFAULT_LIMIT_READ,
) -> NodeList[Node]:
"""`List time series nodes <https://api-docs.cognite.com/20230101/tag/Instances/operation/advancedListInstance>`_.
Expand All @@ -118,6 +138,7 @@ async def list(
space (str | SequenceNotStr[str] | None): Restrict results to this space (or list of spaces).
sort (Sequence[InstanceSort | dict] | InstanceSort | dict | None): Sort order for the results.
filter (Filter | dict[str, Any] | None): Advanced filter to apply. See :class:`~cognite.client.data_classes.filters`.
time_series_type (TimeSeriesType | Sequence[TimeSeriesType] | None): Only return time series of this type (or types). The types are ``"numeric"``, ``"string"`` and ``"state"``. Default: None (all).
limit (int | None): Maximum number of results to return. Defaults to 25. Set to -1, float("inf") or None to return all items.

Returns:
Expand All @@ -135,6 +156,11 @@ async def list(

>>> res = client.data_modeling.time_series.list(space="my-space", limit=None)

List only state time series, or e.g. only numeric and string time series:

>>> res = client.data_modeling.time_series.list(time_series_type="state")
>>> res = client.data_modeling.time_series.list(time_series_type=["numeric", "string"])

Fetch properties from a custom view (note, only time series will be returned), and
apply a custom filter on the name:

Expand All @@ -147,6 +173,7 @@ async def list(
... limit=None,
... )
"""
filter = _build_filter(filter, time_series_type=time_series_type)
sources, strip = resolve_source(source, COGNITE_TIME_SERIES_VIEW_ID)
results = await self._instances_api.list(
instance_type="node",
Expand Down
10 changes: 9 additions & 1 deletion cognite/client/_api/time_series.py
Original file line number Diff line number Diff line change
Expand Up @@ -151,6 +151,10 @@ async def __call__(

Yields:
TimeSeries | TimeSeriesList: yields TimeSeries one by one if chunk_size is not specified, else TimeSeriesList objects.

Note:
State time series are never returned by this method as they are a Data Modeling-only feature (the API leaves them out by default).
You can list these like any other Data Modeling instance, or through the dedicated helper :meth:`client.data_modeling.time_series.list <cognite.client.AsyncCogniteClient.data_modeling.time_series.list>` with ``time_series_type="state"``.
""" # noqa: DOC404
asset_subtree_ids_processed = process_asset_subtree_ids(asset_subtree_ids, asset_subtree_external_ids)
data_set_ids_processed = process_data_set_ids(data_set_ids, data_set_external_ids)
Expand Down Expand Up @@ -784,12 +788,16 @@ async def list(
Returns:
TimeSeriesList: The requested time series.

.. note::
Note:
When using `partitions`, there are few considerations to keep in mind:
* `limit` has to be set to `None` (or `-1`).
* API may reject requests if you specify more than 10 partitions. When Cognite enforces this behavior, the requests result in a 400 Bad Request status.
* Partitions are done independently of sorting: there's no guarantee of the sort order between elements from different partitions. For this reason providing a `sort` parameter when using `partitions` is not allowed.

Note:
State time series are never returned by this method as they are a Data Modeling-only feature (the API leaves them out by default).
You can list these like any other Data Modeling instance, or through the dedicated helper :meth:`client.data_modeling.time_series.list <cognite.client.AsyncCogniteClient.data_modeling.time_series.list>` with ``time_series_type="state"``.

Examples:

List time series:
Expand Down
10 changes: 9 additions & 1 deletion cognite/client/_sync_api/data_modeling/time_series.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

10 changes: 9 additions & 1 deletion cognite/client/_sync_api/time_series.py

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 2 additions & 0 deletions cognite/client/data_classes/time_series.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,8 @@
if TYPE_CHECKING:
from cognite.client.data_classes import Asset, Datapoint

TimeSeriesType: TypeAlias = Literal["numeric", "string", "state"]


class TimeSeries(WriteableCogniteResourceWithClientRef["TimeSeriesWrite"]):
"""This represents a sequence of data points. The TimeSeries object is the metadata about
Expand Down
47 changes: 47 additions & 0 deletions tests/tests_integration/test_api/test_datapoints.py
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
TimeSeries,
TimeSeriesList,
TimeSeriesWrite,
filters,
)
from cognite.client.data_classes.data_modeling import NodeApply, NodeOrEdgeData, Space
from cognite.client.data_classes.data_modeling.cdm.v1 import (
Expand Down Expand Up @@ -494,6 +495,12 @@ def space_for_time_series(cognite_client: CogniteClient) -> Iterator[Space]:
yield cognite_client.data_modeling.spaces.apply(space)


def _dms_ts_listing_name(postfix: str) -> str:
# Shared by a numeric and a state time series, so that we can list both from the classic time series API by name
# (time series created in DM have no classic external ID):
return f"dms-ts-type-listing-long-name-unlikely-to-collide-{postfix}"


@pytest.fixture(scope="session")
def ts_create_in_dms(
cognite_client: CogniteClient, space_for_time_series: Space, os_and_py_version: str
Expand All @@ -504,6 +511,7 @@ def ts_create_in_dms(
external_id=f"dms-time-series-{os_and_py_version}",
is_step=True,
time_series_type="numeric",
name=_dms_ts_listing_name(os_and_py_version),
)
(dms_ts_node,) = cognite_client.data_modeling.instances.apply(dms_ts).nodes
return dms_ts_node
Expand Down Expand Up @@ -567,13 +575,15 @@ def _create_state_time_series(
async_client: AsyncCogniteClient,
space_for_time_series: Space,
state_set: NodeApplyResult,
name: str | None = None,
) -> NodeApplyResult:
state_ts = CogniteTimeSeriesApply(
space=space_for_time_series.space,
external_id=external_id,
is_step=False,
time_series_type="state",
state_set=(state_set.space, state_set.external_id),
name=name,
)
with pytest.MonkeyPatch.context() as mp:
mp.setattr(async_client.data_modeling.instances, "_api_subversion", "beta")
Expand Down Expand Up @@ -615,6 +625,7 @@ def state_ts(
async_client=async_client,
space_for_time_series=space_for_time_series,
state_set=state_set,
name=_dms_ts_listing_name(os_and_py_version),
)


Expand All @@ -635,6 +646,42 @@ def state_ts_b(
)


class TestListStateTimeSeries:
def test_state_time_series_only_listed_from_data_modeling(
self,
cognite_client: CogniteClient,
os_and_py_version: str,
ts_create_in_dms: NodeApplyResult,
state_ts: NodeApplyResult,
) -> None:
# The classic time series API never returns state time series (unless 'includeAllTypes=true' is passed, which
# the SDK doesn't support - and will never tbh as it mixes legacy and DM-only features), -even- when filtering on type:
numeric_id, state_id = ts_create_in_dms.as_id(), state_ts.as_id()
classic = cognite_client.time_series.list(name=_dms_ts_listing_name(os_and_py_version), limit=None)
classic_instance_ids = [ts.instance_id for ts in classic]

assert numeric_id in classic_instance_ids # positive control: the name filter works for DM time series
assert state_id not in classic_instance_ids

# ...but they can of course be listed from the data modeling time series API:
ours = filters.InstanceReferences([numeric_id, state_id])
res = cognite_client.data_modeling.time_series.list(time_series_type="state", filter=ours, limit=None)
assert res.as_ids() == [state_id]

res = cognite_client.data_modeling.time_series.list(
time_series_type=["numeric", "string"], filter=ours, limit=None
)
assert res.as_ids() == [numeric_id]

res = cognite_client.data_modeling.time_series.list(
time_series_type=["numeric", "state"], filter=ours, limit=None
)
assert set(res.as_ids()) == {numeric_id, state_id}

res = cognite_client.data_modeling.time_series.list(filter=ours, limit=None)
assert set(res.as_ids()) == {numeric_id, state_id}


@pytest.mark.allow_no_semaphore(
"StateDatapointsPoster._insert_datapoints holds the semaphore via outer "
"'async with' and calls the http client directly with semaphore=None to avoid double-acquiring."
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,95 @@
from __future__ import annotations

import json
from collections.abc import Sequence
from typing import Any
from unittest.mock import AsyncMock

import pytest

from cognite.client import CogniteClient
from cognite.client._api.data_modeling.time_series import _build_filter
from cognite.client._cognite_client import AsyncCogniteClient
from cognite.client.data_classes.data_modeling import NodeList
from cognite.client.data_classes.filters import Filter
from cognite.client.data_classes.time_series import TimeSeriesType

TYPE_PROPERTY = ["cdf_cdm", "CogniteTimeSeries/v1", "type"]
TYPE_IS_STATE = {"equals": {"property": TYPE_PROPERTY, "value": "state"}}
TYPE_IN_NUMERIC_STRING = {"in": {"property": TYPE_PROPERTY, "values": ["numeric", "string"]}}
SPACE_FILTER = {"equals": {"property": ["node", "space"], "value": "sp"}}


def as_sent(flt: Filter | None) -> dict[str, Any] | None:
# Properties are loaded as tuples of strings, but end up as lists after json has serialized.
# Thus we have this small helper to convert it to the expected format.
if flt is None:
return None
else:
return json.loads(json.dumps(flt.dump(camel_case_property=False)))


class TestBuildFilter:
@pytest.mark.parametrize("filter_as_dict", [True, False])
@pytest.mark.parametrize(
"filter, time_series_type, expected",
[
(None, None, None),
(None, "state", TYPE_IS_STATE), # single -> equals
(None, ["state"], TYPE_IS_STATE), # ...also in a sequence
(None, ("state",), TYPE_IS_STATE),
(None, ["numeric", "string"], TYPE_IN_NUMERIC_STRING), # multiple -> in
(None, ("numeric", "string"), TYPE_IN_NUMERIC_STRING), # any sequence works
(SPACE_FILTER, None, SPACE_FILTER),
(SPACE_FILTER, "state", {"and": [TYPE_IS_STATE, SPACE_FILTER]}),
(SPACE_FILTER, ["numeric", "string"], {"and": [TYPE_IN_NUMERIC_STRING, SPACE_FILTER]}),
],
)
def test_build_filter(
self,
filter: dict[str, Any] | None,
time_series_type: TimeSeriesType | Sequence[TimeSeriesType] | None,
expected: dict[str, Any] | None,
filter_as_dict: bool,
) -> None:
Comment thread
haakonvt marked this conversation as resolved.
given: Filter | dict | None = filter
if not (filter is None or filter_as_dict):
given = Filter.load(filter)

assert expected == as_sent(_build_filter(given, time_series_type=time_series_type))

@pytest.mark.parametrize("time_series_type", [[], ()])
def test_build_filter_raises_on_empty_time_series_types(self, time_series_type: Sequence[TimeSeriesType]) -> None:
with pytest.raises(ValueError, match="'time_series_type' must not be empty, pass None"):
_build_filter(None, time_series_type=time_series_type)


class TestDMTimeSeriesListTimeSeriesTypes:
@pytest.fixture
def list_mock(self, async_client: AsyncCogniteClient, monkeypatch: pytest.MonkeyPatch) -> AsyncMock:
mock = AsyncMock(return_value=NodeList([]))
monkeypatch.setattr(async_client.data_modeling.instances, "list", mock)
return mock

@staticmethod
def sent_filter(list_mock: AsyncMock) -> dict[str, Any] | None:
flt = list_mock.call_args.kwargs["filter"]
if flt is None or isinstance(flt, dict):
return flt
return json.loads(json.dumps(flt.dump(camel_case_property=False)))

@pytest.mark.parametrize(
"kwargs, expected",
[
({}, None),
({"time_series_type": "state"}, TYPE_IS_STATE),
({"time_series_type": ["numeric", "string"]}, TYPE_IN_NUMERIC_STRING),
({"filter": SPACE_FILTER}, SPACE_FILTER),
({"time_series_type": "state", "filter": SPACE_FILTER}, {"and": [TYPE_IS_STATE, SPACE_FILTER]}),
],
)
def test_time_series_types_is_combined_with_filter(
self, cognite_client: CogniteClient, list_mock: AsyncMock, kwargs: dict[str, Any], expected: dict | None
) -> None:
cognite_client.data_modeling.time_series.list(**kwargs)
assert self.sent_filter(list_mock) == expected
Loading