From f8ff4c21152d4eac221a0ca1a883d6e2b1b41de0 Mon Sep 17 00:00:00 2001 From: Nathan Thorpe Date: Mon, 28 Sep 2026 15:03:04 -0700 Subject: [PATCH] add task lineage endpoint --- .../v1/api/execution/get_task_lineage.py | 194 ++++++++++++++++++ cirro_api_client/v1/models/__init__.py | 6 +- .../v1/models/get_task_lineage_response.py | 74 +++++++ ...y => get_task_lineage_response_lineage.py} | 12 +- .../v1/models/sheet_query_response.py | 17 +- cirro_api_client/v1/models/task.py | 20 ++ pyproject.toml | 2 +- 7 files changed, 309 insertions(+), 16 deletions(-) create mode 100644 cirro_api_client/v1/api/execution/get_task_lineage.py create mode 100644 cirro_api_client/v1/models/get_task_lineage_response.py rename cirro_api_client/v1/models/{sheet_query_response_rows_item.py => get_task_lineage_response_lineage.py} (77%) diff --git a/cirro_api_client/v1/api/execution/get_task_lineage.py b/cirro_api_client/v1/api/execution/get_task_lineage.py new file mode 100644 index 0000000..5b83dff --- /dev/null +++ b/cirro_api_client/v1/api/execution/get_task_lineage.py @@ -0,0 +1,194 @@ +from http import HTTPStatus +from typing import Any +from urllib.parse import quote + +import httpx + +from ... import errors +from ...client import Client +from ...models.get_task_lineage_response import GetTaskLineageResponse +from ...types import Response + + +def _get_kwargs( + project_id: str, + dataset_id: str, + task_id: str, +) -> dict[str, Any]: + _kwargs: dict[str, Any] = { + "method": "get", + "url": "/projects/{project_id}/execution/{dataset_id}/tasks/{task_id}/lineage".format( + project_id=quote(str(project_id), safe=""), + dataset_id=quote(str(dataset_id), safe=""), + task_id=quote(str(task_id), safe=""), + ), + } + + return _kwargs + + +def _parse_response(*, client: Client, response: httpx.Response) -> GetTaskLineageResponse | None: + if response.status_code == 200: + response_200 = GetTaskLineageResponse.from_dict(response.json()) + + return response_200 + + errors.handle_error_response(response, client.raise_on_unexpected_status) + + +def _build_response(*, client: Client, response: httpx.Response) -> Response[GetTaskLineageResponse]: + return Response( + status_code=HTTPStatus(response.status_code), + content=response.content, + headers=response.headers, + parsed=_parse_response(client=client, response=response), + ) + + +def sync_detailed( + project_id: str, + dataset_id: str, + task_id: str, + *, + client: Client, +) -> Response[GetTaskLineageResponse]: + """Get task lineage + + Gets detailed lineage information on the individual Nextflow task + + Args: + project_id (str): + dataset_id (str): + task_id (str): + client (Client): instance of the API client + + Raises: + errors.UnexpectedStatus: If the server returns an undocumented status code and Client.raise_on_unexpected_status is True. + httpx.TimeoutException: If the request takes longer than Client.timeout. + + Returns: + Response[GetTaskLineageResponse] + """ + + kwargs = _get_kwargs( + project_id=project_id, + dataset_id=dataset_id, + task_id=task_id, + ) + + response = client.get_httpx_client().request( + auth=client.get_auth(), + **kwargs, + ) + + return _build_response(client=client, response=response) + + +def sync( + project_id: str, + dataset_id: str, + task_id: str, + *, + client: Client, +) -> GetTaskLineageResponse | None: + """Get task lineage + + Gets detailed lineage information on the individual Nextflow task + + Args: + project_id (str): + dataset_id (str): + task_id (str): + client (Client): instance of the API client + + Raises: + errors.UnexpectedStatus: If the server returns an undocumented status code and Client.raise_on_unexpected_status is True. + httpx.TimeoutException: If the request takes longer than Client.timeout. + + Returns: + GetTaskLineageResponse + """ + + try: + return sync_detailed( + project_id=project_id, + dataset_id=dataset_id, + task_id=task_id, + client=client, + ).parsed + except errors.NotFoundException: + return None + + +async def asyncio_detailed( + project_id: str, + dataset_id: str, + task_id: str, + *, + client: Client, +) -> Response[GetTaskLineageResponse]: + """Get task lineage + + Gets detailed lineage information on the individual Nextflow task + + Args: + project_id (str): + dataset_id (str): + task_id (str): + client (Client): instance of the API client + + Raises: + errors.UnexpectedStatus: If the server returns an undocumented status code and Client.raise_on_unexpected_status is True. + httpx.TimeoutException: If the request takes longer than Client.timeout. + + Returns: + Response[GetTaskLineageResponse] + """ + + kwargs = _get_kwargs( + project_id=project_id, + dataset_id=dataset_id, + task_id=task_id, + ) + + response = await client.get_async_httpx_client().request(auth=client.get_auth(), **kwargs) + + return _build_response(client=client, response=response) + + +async def asyncio( + project_id: str, + dataset_id: str, + task_id: str, + *, + client: Client, +) -> GetTaskLineageResponse | None: + """Get task lineage + + Gets detailed lineage information on the individual Nextflow task + + Args: + project_id (str): + dataset_id (str): + task_id (str): + client (Client): instance of the API client + + Raises: + errors.UnexpectedStatus: If the server returns an undocumented status code and Client.raise_on_unexpected_status is True. + httpx.TimeoutException: If the request takes longer than Client.timeout. + + Returns: + GetTaskLineageResponse + """ + + try: + return ( + await asyncio_detailed( + project_id=project_id, + dataset_id=dataset_id, + task_id=task_id, + client=client, + ) + ).parsed + except errors.NotFoundException: + return None diff --git a/cirro_api_client/v1/models/__init__.py b/cirro_api_client/v1/models/__init__.py index c0d4b2a..d213296 100644 --- a/cirro_api_client/v1/models/__init__.py +++ b/cirro_api_client/v1/models/__init__.py @@ -106,6 +106,8 @@ from .get_execution_logs_response import GetExecutionLogsResponse from .get_project_summary_response_200 import GetProjectSummaryResponse200 from .get_task_files_response import GetTaskFilesResponse +from .get_task_lineage_response import GetTaskLineageResponse +from .get_task_lineage_response_lineage import GetTaskLineageResponseLineage from .governance_access_type import GovernanceAccessType from .governance_classification import GovernanceClassification from .governance_contact import GovernanceContact @@ -237,7 +239,6 @@ from .sheet_job_type import SheetJobType from .sheet_query_request import SheetQueryRequest from .sheet_query_response import SheetQueryResponse -from .sheet_query_response_rows_item import SheetQueryResponseRowsItem from .sheet_sort import SheetSort from .sheet_type import SheetType from .sheet_update_response import SheetUpdateResponse @@ -391,6 +392,8 @@ "GetExecutionLogsResponse", "GetProjectSummaryResponse200", "GetTaskFilesResponse", + "GetTaskLineageResponse", + "GetTaskLineageResponseLineage", "GovernanceAccessType", "GovernanceClassification", "GovernanceContact", @@ -522,7 +525,6 @@ "SheetJobType", "SheetQueryRequest", "SheetQueryResponse", - "SheetQueryResponseRowsItem", "SheetSort", "SheetType", "SheetUpdateResponse", diff --git a/cirro_api_client/v1/models/get_task_lineage_response.py b/cirro_api_client/v1/models/get_task_lineage_response.py new file mode 100644 index 0000000..badfcb9 --- /dev/null +++ b/cirro_api_client/v1/models/get_task_lineage_response.py @@ -0,0 +1,74 @@ +from __future__ import annotations + +from collections.abc import Mapping +from typing import TYPE_CHECKING, Any, TypeVar + +from attrs import define as _attrs_define +from attrs import field as _attrs_field + +from ..types import UNSET, Unset + +if TYPE_CHECKING: + from ..models.get_task_lineage_response_lineage import GetTaskLineageResponseLineage + + +T = TypeVar("T", bound="GetTaskLineageResponse") + + +@_attrs_define +class GetTaskLineageResponse: + """ + Attributes: + lineage (GetTaskLineageResponseLineage | Unset): Lineage specific to the executor + """ + + lineage: GetTaskLineageResponseLineage | Unset = UNSET + additional_properties: dict[str, Any] = _attrs_field(init=False, factory=dict) + + def to_dict(self) -> dict[str, Any]: + lineage: dict[str, Any] | Unset = UNSET + if not isinstance(self.lineage, Unset): + lineage = self.lineage.to_dict() + + field_dict: dict[str, Any] = {} + field_dict.update(self.additional_properties) + field_dict.update({}) + if lineage is not UNSET: + field_dict["lineage"] = lineage + + return field_dict + + @classmethod + def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: + from ..models.get_task_lineage_response_lineage import GetTaskLineageResponseLineage + + d = dict(src_dict) + _lineage = d.pop("lineage", UNSET) + lineage: GetTaskLineageResponseLineage | Unset + if isinstance(_lineage, Unset): + lineage = UNSET + else: + lineage = GetTaskLineageResponseLineage.from_dict(_lineage) + + get_task_lineage_response = cls( + lineage=lineage, + ) + + get_task_lineage_response.additional_properties = d + return get_task_lineage_response + + @property + def additional_keys(self) -> list[str]: + return list(self.additional_properties.keys()) + + def __getitem__(self, key: str) -> Any: + return self.additional_properties[key] + + def __setitem__(self, key: str, value: Any) -> None: + self.additional_properties[key] = value + + def __delitem__(self, key: str) -> None: + del self.additional_properties[key] + + def __contains__(self, key: str) -> bool: + return key in self.additional_properties diff --git a/cirro_api_client/v1/models/sheet_query_response_rows_item.py b/cirro_api_client/v1/models/get_task_lineage_response_lineage.py similarity index 77% rename from cirro_api_client/v1/models/sheet_query_response_rows_item.py rename to cirro_api_client/v1/models/get_task_lineage_response_lineage.py index 12dae57..c683a52 100644 --- a/cirro_api_client/v1/models/sheet_query_response_rows_item.py +++ b/cirro_api_client/v1/models/get_task_lineage_response_lineage.py @@ -6,12 +6,12 @@ from attrs import define as _attrs_define from attrs import field as _attrs_field -T = TypeVar("T", bound="SheetQueryResponseRowsItem") +T = TypeVar("T", bound="GetTaskLineageResponseLineage") @_attrs_define -class SheetQueryResponseRowsItem: - """ """ +class GetTaskLineageResponseLineage: + """Lineage specific to the executor""" additional_properties: dict[str, Any] = _attrs_field(init=False, factory=dict) @@ -25,10 +25,10 @@ def to_dict(self) -> dict[str, Any]: @classmethod def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: d = dict(src_dict) - sheet_query_response_rows_item = cls() + get_task_lineage_response_lineage = cls() - sheet_query_response_rows_item.additional_properties = d - return sheet_query_response_rows_item + get_task_lineage_response_lineage.additional_properties = d + return get_task_lineage_response_lineage @property def additional_keys(self) -> list[str]: diff --git a/cirro_api_client/v1/models/sheet_query_response.py b/cirro_api_client/v1/models/sheet_query_response.py index 50ebcf8..f36f2cd 100644 --- a/cirro_api_client/v1/models/sheet_query_response.py +++ b/cirro_api_client/v1/models/sheet_query_response.py @@ -1,14 +1,13 @@ from __future__ import annotations from collections.abc import Mapping -from typing import TYPE_CHECKING, Any, TypeVar +from typing import TYPE_CHECKING, Any, TypeVar, cast from attrs import define as _attrs_define from attrs import field as _attrs_field if TYPE_CHECKING: from ..models.query_column import QueryColumn - from ..models.sheet_query_response_rows_item import SheetQueryResponseRowsItem T = TypeVar("T", bound="SheetQueryResponse") @@ -22,12 +21,12 @@ class SheetQueryResponse: Attributes: columns (list[QueryColumn]): column definitions, starting with `_row_id` - rows (list[list[SheetQueryResponseRowsItem]]): row data, each list aligned with `columns` + rows (list[list[bool | float | int | str]]): row data, each list aligned with `columns` total_row_count (int): number of total rows in the result set """ columns: list[QueryColumn] - rows: list[list[SheetQueryResponseRowsItem]] + rows: list[list[bool | float | int | str]] total_row_count: int additional_properties: dict[str, Any] = _attrs_field(init=False, factory=dict) @@ -41,7 +40,8 @@ def to_dict(self) -> dict[str, Any]: for rows_item_data in self.rows: rows_item = [] for rows_item_item_data in rows_item_data: - rows_item_item = rows_item_item_data.to_dict() + rows_item_item: bool | float | int | str + rows_item_item = rows_item_item_data rows_item.append(rows_item_item) rows.append(rows_item) @@ -63,7 +63,6 @@ def to_dict(self) -> dict[str, Any]: @classmethod def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: from ..models.query_column import QueryColumn - from ..models.sheet_query_response_rows_item import SheetQueryResponseRowsItem d = dict(src_dict) columns = [] @@ -79,7 +78,11 @@ def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: rows_item = [] _rows_item = rows_item_data for rows_item_item_data in _rows_item: - rows_item_item = SheetQueryResponseRowsItem.from_dict(rows_item_item_data) + + def _parse_rows_item_item(data: object) -> bool | float | int | str: + return cast(bool | float | int | str, data) + + rows_item_item = _parse_rows_item_item(rows_item_item_data) rows_item.append(rows_item_item) diff --git a/cirro_api_client/v1/models/task.py b/cirro_api_client/v1/models/task.py index b690866..baaa1b4 100644 --- a/cirro_api_client/v1/models/task.py +++ b/cirro_api_client/v1/models/task.py @@ -18,6 +18,7 @@ class Task: Attributes: name (str): status (str): + hash_ (None | str | Unset): Hash of task from executor (used for caching) native_job_id (None | str | Unset): Job ID on the underlying execution environment (i.e. AWS Batch ID) status_message (None | str | Unset): requested_at (datetime.datetime | None | Unset): @@ -32,6 +33,7 @@ class Task: name: str status: str + hash_: None | str | Unset = UNSET native_job_id: None | str | Unset = UNSET status_message: None | str | Unset = UNSET requested_at: datetime.datetime | None | Unset = UNSET @@ -49,6 +51,12 @@ def to_dict(self) -> dict[str, Any]: status = self.status + hash_: None | str | Unset + if isinstance(self.hash_, Unset): + hash_ = UNSET + else: + hash_ = self.hash_ + native_job_id: None | str | Unset if isinstance(self.native_job_id, Unset): native_job_id = UNSET @@ -123,6 +131,8 @@ def to_dict(self) -> dict[str, Any]: "status": status, } ) + if hash_ is not UNSET: + field_dict["hash"] = hash_ if native_job_id is not UNSET: field_dict["nativeJobId"] = native_job_id if status_message is not UNSET: @@ -153,6 +163,15 @@ def from_dict(cls: type[T], src_dict: Mapping[str, Any]) -> T: status = d.pop("status") + def _parse_hash_(data: object) -> None | str | Unset: + if data is None: + return data + if isinstance(data, Unset): + return data + return cast(None | str | Unset, data) + + hash_ = _parse_hash_(d.pop("hash", UNSET)) + def _parse_native_job_id(data: object) -> None | str | Unset: if data is None: return data @@ -270,6 +289,7 @@ def _parse_log_location(data: object) -> None | str | Unset: task = cls( name=name, status=status, + hash_=hash_, native_job_id=native_job_id, status_message=status_message, requested_at=requested_at, diff --git a/pyproject.toml b/pyproject.toml index 6cfb528..a0d2a24 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [tool.poetry] name = "cirro_api_client" -version = "1.7.0" +version = "1.7.1" description = "A client library for accessing Cirro" authors = ["Cirro "] license = "MIT"