From 4067c8b80b6983a601fd314106acb69234eb88ce Mon Sep 17 00:00:00 2001 From: Alex Fieraru Date: Fri, 8 May 2026 15:24:19 +0300 Subject: [PATCH 1/4] feat: create fetch_data() function --- dataconnect/client.py | 12 ++++++++ dataconnect/service/base.py | 8 ++++++ dataconnect/service/default.py | 28 ++++++++++++++++++- dataconnect/service/mappers.py | 12 ++++++++ .../transport/arrow_flight/transport.py | 1 + 5 files changed, 60 insertions(+), 1 deletion(-) diff --git a/dataconnect/client.py b/dataconnect/client.py index 84ce20d..a2a87d6 100644 --- a/dataconnect/client.py +++ b/dataconnect/client.py @@ -12,6 +12,10 @@ from dataconnect.models import Study from dataconnect.service import DataConnectService, DefaultDataConnectService +from uuid import UUID + +import pandas as pd + _DEFAULT_HOST = "enodia-gateway.platform.imedidata.com" _DEFAULT_PORT = 443 @@ -45,6 +49,14 @@ def connect( def get_studies(self) -> list[Study]: """List the studies the client is authorized to access.""" return self._service.get_studies() + + def fetch_data( + self, + dataset_uuid: UUID, + first_n_rows: int | None = None, + ) -> pd.DataFrame: + """Fetch data for a given dataset UUID.""" + return self._service.fetch_data(dataset_uuid, first_n_rows) # Lifecycle diff --git a/dataconnect/service/base.py b/dataconnect/service/base.py index 78a26b0..213b2e2 100644 --- a/dataconnect/service/base.py +++ b/dataconnect/service/base.py @@ -6,6 +6,7 @@ from dataconnect.models import Study +import pandas as pd class DataConnectService(ABC): """Abstract service interface — defines all operations available to the client.""" @@ -13,5 +14,12 @@ class DataConnectService(ABC): @abstractmethod def get_studies(self) -> list[Study]: ... + @abstractmethod + def fetch_data( + self, + dataset_uuid: str, + first_n_rows: int | None = None, + ) -> pd.DataFrame: ... + @abstractmethod def close(self) -> None: ... diff --git a/dataconnect/service/default.py b/dataconnect/service/default.py index dd287cf..607fbde 100644 --- a/dataconnect/service/default.py +++ b/dataconnect/service/default.py @@ -14,7 +14,7 @@ ) from dataconnect.models import Study from dataconnect.service.base import DataConnectService -from dataconnect.service.mappers import resource_to_study +from dataconnect.service.mappers import resource_to_fetched_data, resource_to_study from dataconnect.transport.base import Transport from dataconnect.transport.errors import ( TransportAuthenticationError, @@ -26,9 +26,11 @@ TransportStatusError, ) from dataconnect.transport.models import ResourceQuery +from uuid import UUID # Server action identifiers _ACTION_LIST_STUDIES = "studies.list" +_ACTION_FETCH_TICKET = "data.fetch_ticket" def _translate_error(ex: TransportError) -> DataConnectError: @@ -72,6 +74,30 @@ def get_studies(self) -> list[Study]: except (IndexError, KeyError, TypeError, ValueError) as ex: raise ValidationError(f"Unexpected studies response format: {ex}") from ex + def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None): + + if not dataset_uuid or not str(dataset_uuid).strip(): + raise ValueError("dataset_uuid must be provided.") + + if first_n_rows is not None and first_n_rows <= 0: + raise ValueError("first_n_rows must be a positive integer when provided.") + + request = ResourceQuery(action=_ACTION_FETCH_TICKET).append_body( + "dataset_uuid", str(dataset_uuid), + limit=first_n_rows, + ) + + try: + resources = self._transport.list_resources(request) + except TransportError as ex: + raise _translate_error(ex) from ex + + try: + return [resource_to_fetched_data(r) for r in resources] + except (IndexError, KeyError, TypeError, ValueError) as ex: + raise ValidationError(f"Unexpected data fetch response format: {ex}") from ex + + def close(self) -> None: try: diff --git a/dataconnect/service/mappers.py b/dataconnect/service/mappers.py index a4a6407..9261208 100644 --- a/dataconnect/service/mappers.py +++ b/dataconnect/service/mappers.py @@ -8,6 +8,9 @@ from __future__ import annotations import json +import pandas as pd + +from turtle import pd from uuid import UUID from dataconnect.exceptions import NotFoundError @@ -28,3 +31,12 @@ def resource_to_study(resource: ResourceInfo) -> Study: name=data["name"], environments=[StudyEnvironment(uuid=UUID(e["uuid"]), name=e["name"]) for e in data.get("environments", [])], ) + +def resource_to_fetched_data(resource: ResourceInfo) -> pd.DataFrame: + """Parse a transport-layer ``ResourceInfo`` into pandas DataFrame.""" + + if not resource or not resource.endpoints or not resource.endpoints[0].ticket: + raise NotFoundError("Invalid resource: missing endpoints or ticket") + + data = json.loads(resource.endpoints[0].ticket.decode("utf-8")) + return data.to_pandas() diff --git a/dataconnect/transport/arrow_flight/transport.py b/dataconnect/transport/arrow_flight/transport.py index cd714ee..7e86b8a 100644 --- a/dataconnect/transport/arrow_flight/transport.py +++ b/dataconnect/transport/arrow_flight/transport.py @@ -39,6 +39,7 @@ def _to_resource_info(info: flight.FlightInfo) -> ResourceInfo: # Maps service-layer action names to the flight_type value the Arrow Flight server expects. _ACTION_FLIGHT_TYPE: dict[str, str] = { "studies.list": "STUDIES", + "data.fetch_ticket": "DATA_FETCH_TICKET", } From 1681644eb4cddae998ce6700b81b3001bbcec901 Mon Sep 17 00:00:00 2001 From: Alex Fieraru Date: Fri, 8 May 2026 15:46:23 +0300 Subject: [PATCH 2/4] feat: fix copilot recommendations --- dataconnect/client.py | 4 ++-- dataconnect/service/base.py | 3 ++- dataconnect/service/default.py | 17 +++++++++++++---- dataconnect/service/mappers.py | 16 +++++++++++----- pyproject.toml | 2 +- 5 files changed, 29 insertions(+), 13 deletions(-) diff --git a/dataconnect/client.py b/dataconnect/client.py index a2a87d6..820716b 100644 --- a/dataconnect/client.py +++ b/dataconnect/client.py @@ -54,8 +54,8 @@ def fetch_data( self, dataset_uuid: UUID, first_n_rows: int | None = None, - ) -> pd.DataFrame: - """Fetch data for a given dataset UUID.""" + ) -> pd.DataFrame: + """Fetch data frames for a given dataset UUID.""" return self._service.fetch_data(dataset_uuid, first_n_rows) # Lifecycle diff --git a/dataconnect/service/base.py b/dataconnect/service/base.py index 213b2e2..d48ed04 100644 --- a/dataconnect/service/base.py +++ b/dataconnect/service/base.py @@ -3,6 +3,7 @@ from __future__ import annotations from abc import ABC, abstractmethod +from uuid import UUID from dataconnect.models import Study @@ -17,7 +18,7 @@ def get_studies(self) -> list[Study]: ... @abstractmethod def fetch_data( self, - dataset_uuid: str, + dataset_uuid: UUID, first_n_rows: int | None = None, ) -> pd.DataFrame: ... diff --git a/dataconnect/service/default.py b/dataconnect/service/default.py index 607fbde..80d626d 100644 --- a/dataconnect/service/default.py +++ b/dataconnect/service/default.py @@ -28,6 +28,8 @@ from dataconnect.transport.models import ResourceQuery from uuid import UUID +import pandas as pd + # Server action identifiers _ACTION_LIST_STUDIES = "studies.list" _ACTION_FETCH_TICKET = "data.fetch_ticket" @@ -74,7 +76,7 @@ def get_studies(self) -> list[Study]: except (IndexError, KeyError, TypeError, ValueError) as ex: raise ValidationError(f"Unexpected studies response format: {ex}") from ex - def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None): + def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None) -> pd.DataFrame: if not dataset_uuid or not str(dataset_uuid).strip(): raise ValueError("dataset_uuid must be provided.") @@ -83,8 +85,10 @@ def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None): raise ValueError("first_n_rows must be a positive integer when provided.") request = ResourceQuery(action=_ACTION_FETCH_TICKET).append_body( - "dataset_uuid", str(dataset_uuid), - limit=first_n_rows, + { + "dataset_uuid": str(dataset_uuid), + "limit": first_n_rows, + } ) try: @@ -93,10 +97,15 @@ def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None): raise _translate_error(ex) from ex try: - return [resource_to_fetched_data(r) for r in resources] + frames = [resource_to_fetched_data(r) for r in resources] except (IndexError, KeyError, TypeError, ValueError) as ex: raise ValidationError(f"Unexpected data fetch response format: {ex}") from ex + if not frames: + return pd.DataFrame() + + return pd.concat(frames, ignore_index=True) + def close(self) -> None: diff --git a/dataconnect/service/mappers.py b/dataconnect/service/mappers.py index 9261208..0e36f08 100644 --- a/dataconnect/service/mappers.py +++ b/dataconnect/service/mappers.py @@ -9,8 +9,8 @@ import json import pandas as pd +import pyarrow as pa -from turtle import pd from uuid import UUID from dataconnect.exceptions import NotFoundError @@ -33,10 +33,16 @@ def resource_to_study(resource: ResourceInfo) -> Study: ) def resource_to_fetched_data(resource: ResourceInfo) -> pd.DataFrame: - """Parse a transport-layer ``ResourceInfo`` into pandas DataFrame.""" + """Parse a transport-layer ``ResourceInfo`` into pandas DataFrame. + + The ticket bytes are expected to contain an Arrow IPC stream. They are + read via ``pyarrow.ipc.open_stream`` and converted to a ``pd.DataFrame``. + """ if not resource or not resource.endpoints or not resource.endpoints[0].ticket: raise NotFoundError("Invalid resource: missing endpoints or ticket") - - data = json.loads(resource.endpoints[0].ticket.decode("utf-8")) - return data.to_pandas() + + ticket_bytes = resource.endpoints[0].ticket + buf = pa.BufferReader(ticket_bytes) + reader = pa.ipc.open_stream(buf) + return reader.read_all().to_pandas() diff --git a/pyproject.toml b/pyproject.toml index a0d5680..61f165d 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -19,6 +19,7 @@ url = "https://mdsol.jfrog.io/artifactory/api/pypi/pypi-prod-virtual/simple" [tool.poetry.dependencies] python = "^3.13" pyarrow = "^19.0.0" +pandas = "^2.0.2" # SERVICE [tool.poetry.group.service] @@ -32,7 +33,6 @@ gunicorn ="^20.1.0" optional = true [tool.poetry.group.ml.dependencies] -pandas = "^2.0.2" # DEV [tool.poetry.group.dev] From c216d2f548e471d0fda6bbd49c6be82b59d8398d Mon Sep 17 00:00:00 2001 From: Alex Fieraru Date: Fri, 8 May 2026 16:34:14 +0300 Subject: [PATCH 3/4] feat: fix linting --- dataconnect/client.py | 9 ++++----- dataconnect/service/base.py | 3 ++- dataconnect/service/default.py | 10 +++++----- dataconnect/service/mappers.py | 5 +++-- 4 files changed, 14 insertions(+), 13 deletions(-) diff --git a/dataconnect/client.py b/dataconnect/client.py index 820716b..83a0258 100644 --- a/dataconnect/client.py +++ b/dataconnect/client.py @@ -8,14 +8,13 @@ from __future__ import annotations from types import TracebackType - -from dataconnect.models import Study -from dataconnect.service import DataConnectService, DefaultDataConnectService - from uuid import UUID import pandas as pd +from dataconnect.models import Study +from dataconnect.service import DataConnectService, DefaultDataConnectService + _DEFAULT_HOST = "enodia-gateway.platform.imedidata.com" _DEFAULT_PORT = 443 @@ -49,7 +48,7 @@ def connect( def get_studies(self) -> list[Study]: """List the studies the client is authorized to access.""" return self._service.get_studies() - + def fetch_data( self, dataset_uuid: UUID, diff --git a/dataconnect/service/base.py b/dataconnect/service/base.py index d48ed04..da54bab 100644 --- a/dataconnect/service/base.py +++ b/dataconnect/service/base.py @@ -5,9 +5,10 @@ from abc import ABC, abstractmethod from uuid import UUID +import pandas as pd + from dataconnect.models import Study -import pandas as pd class DataConnectService(ABC): """Abstract service interface — defines all operations available to the client.""" diff --git a/dataconnect/service/default.py b/dataconnect/service/default.py index 80d626d..20d2fd4 100644 --- a/dataconnect/service/default.py +++ b/dataconnect/service/default.py @@ -2,6 +2,10 @@ from __future__ import annotations +from uuid import UUID + +import pandas as pd + from dataconnect.exceptions import ( AuthenticationError, AuthorizationError, @@ -26,9 +30,6 @@ TransportStatusError, ) from dataconnect.transport.models import ResourceQuery -from uuid import UUID - -import pandas as pd # Server action identifiers _ACTION_LIST_STUDIES = "studies.list" @@ -95,7 +96,7 @@ def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None) -> pd. resources = self._transport.list_resources(request) except TransportError as ex: raise _translate_error(ex) from ex - + try: frames = [resource_to_fetched_data(r) for r in resources] except (IndexError, KeyError, TypeError, ValueError) as ex: @@ -106,7 +107,6 @@ def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None) -> pd. return pd.concat(frames, ignore_index=True) - def close(self) -> None: try: diff --git a/dataconnect/service/mappers.py b/dataconnect/service/mappers.py index 0e36f08..3f7346c 100644 --- a/dataconnect/service/mappers.py +++ b/dataconnect/service/mappers.py @@ -8,11 +8,11 @@ from __future__ import annotations import json +from uuid import UUID + import pandas as pd import pyarrow as pa -from uuid import UUID - from dataconnect.exceptions import NotFoundError from dataconnect.models import Study, StudyEnvironment from dataconnect.transport.models import ResourceInfo @@ -32,6 +32,7 @@ def resource_to_study(resource: ResourceInfo) -> Study: environments=[StudyEnvironment(uuid=UUID(e["uuid"]), name=e["name"]) for e in data.get("environments", [])], ) + def resource_to_fetched_data(resource: ResourceInfo) -> pd.DataFrame: """Parse a transport-layer ``ResourceInfo`` into pandas DataFrame. From 67f72a0cccebbaabbbaef8b68c7f8e21bfc5af60 Mon Sep 17 00:00:00 2001 From: Alex Fieraru Date: Mon, 11 May 2026 15:30:27 +0300 Subject: [PATCH 4/4] feat: rework after code refactorization --- dataconnect/service/default.py | 17 ++---- dataconnect/service/mappers.py | 20 +++--- .../transport/arrow_flight/transport.py | 61 ++++++++++++++++++- dataconnect/transport/base.py | 10 ++- dataconnect/transport/models.py | 10 +++ 5 files changed, 91 insertions(+), 27 deletions(-) diff --git a/dataconnect/service/default.py b/dataconnect/service/default.py index 20d2fd4..8a1fc7b 100644 --- a/dataconnect/service/default.py +++ b/dataconnect/service/default.py @@ -18,7 +18,7 @@ ) from dataconnect.models import Study from dataconnect.service.base import DataConnectService -from dataconnect.service.mappers import resource_to_fetched_data, resource_to_study +from dataconnect.service.mappers import resource_to_study, resource_to_fetched_data from dataconnect.transport.base import Transport from dataconnect.transport.errors import ( TransportAuthenticationError, @@ -87,26 +87,19 @@ def fetch_data(self, dataset_uuid: UUID, first_n_rows: int | None = None) -> pd. request = ResourceQuery(action=_ACTION_FETCH_TICKET).append_body( { + "study_env_uuid": None, + "dataset_name": None, "dataset_uuid": str(dataset_uuid), "limit": first_n_rows, } ) try: - resources = self._transport.list_resources(request) + table = self._transport.do_get(request) + return resource_to_fetched_data(table) except TransportError as ex: raise _translate_error(ex) from ex - try: - frames = [resource_to_fetched_data(r) for r in resources] - except (IndexError, KeyError, TypeError, ValueError) as ex: - raise ValidationError(f"Unexpected data fetch response format: {ex}") from ex - - if not frames: - return pd.DataFrame() - - return pd.concat(frames, ignore_index=True) - def close(self) -> None: try: diff --git a/dataconnect/service/mappers.py b/dataconnect/service/mappers.py index 3f7346c..807cd9e 100644 --- a/dataconnect/service/mappers.py +++ b/dataconnect/service/mappers.py @@ -15,7 +15,7 @@ from dataconnect.exceptions import NotFoundError from dataconnect.models import Study, StudyEnvironment -from dataconnect.transport.models import ResourceInfo +from dataconnect.transport.models import ResourceInfo, DataTable def resource_to_study(resource: ResourceInfo) -> Study: @@ -33,17 +33,11 @@ def resource_to_study(resource: ResourceInfo) -> Study: ) -def resource_to_fetched_data(resource: ResourceInfo) -> pd.DataFrame: - """Parse a transport-layer ``ResourceInfo`` into pandas DataFrame. +def resource_to_fetched_data(table: DataTable) -> pd.DataFrame: + """Convert a transport-layer ``DataTable`` into a ``pandas.DataFrame``.""" - The ticket bytes are expected to contain an Arrow IPC stream. They are - read via ``pyarrow.ipc.open_stream`` and converted to a ``pd.DataFrame``. - """ + ipc_buffer = pa.BufferReader(table.ipc_bytes) + reader = pa.ipc.open_stream(ipc_buffer) + table = reader.read_all() - if not resource or not resource.endpoints or not resource.endpoints[0].ticket: - raise NotFoundError("Invalid resource: missing endpoints or ticket") - - ticket_bytes = resource.endpoints[0].ticket - buf = pa.BufferReader(ticket_bytes) - reader = pa.ipc.open_stream(buf) - return reader.read_all().to_pandas() + return pd.DataFrame(table.to_pandas()) \ No newline at end of file diff --git a/dataconnect/transport/arrow_flight/transport.py b/dataconnect/transport/arrow_flight/transport.py index 7e86b8a..680e39f 100644 --- a/dataconnect/transport/arrow_flight/transport.py +++ b/dataconnect/transport/arrow_flight/transport.py @@ -9,6 +9,7 @@ import json +import pyarrow as pa from pyarrow import flight from dataconnect.transport.base import Transport @@ -18,7 +19,7 @@ TransportConnectionError, TransportStatusError, ) -from dataconnect.transport.models import DataRef, ResourceInfo, ResourceQuery +from dataconnect.transport.models import DataRef, ResourceInfo, ResourceQuery, DataTable def _to_resource_info(info: flight.FlightInfo) -> ResourceInfo: @@ -35,6 +36,23 @@ def _to_resource_info(info: flight.FlightInfo) -> ResourceInfo: total_records=info.total_records, ) +def _to_bytes(table: pa.Table) -> DataTable: + """Serialize a ``pa.Table`` to a technology-agnostic ``DataTable``. + + Each record batch is serialized individually as Arrow IPC bytes. + The schema is serialized separately so it can be recovered without + the data batches. + """ + schema_bytes = table.schema.serialize().to_pybytes() + + sink = pa.BufferOutputStream() + writer = pa.ipc.new_stream(sink, table.schema) + for batch in table.to_batches(): + writer.write_batch(batch) + writer.close() + ipc_bytes = sink.getvalue().to_pybytes() + + return DataTable(schema_bytes=schema_bytes, ipc_bytes=ipc_bytes) # Maps service-layer action names to the flight_type value the Arrow Flight server expects. _ACTION_FLIGHT_TYPE: dict[str, str] = { @@ -102,5 +120,46 @@ def list_resources(self, request: ResourceQuery) -> list[ResourceInfo]: except Exception as ex: raise TransportConnectionError(f"Unexpected error during list_resources: {ex}") from ex + def do_get(self, request: ResourceQuery) -> DataTable: + """Call FlightClient.do_get and read all chunks into a single pa.Table.""" + + flight_type = _ACTION_FLIGHT_TYPE.get(request.action) + + if flight_type is None: + raise TransportStatusError( + f"Unknown action: {request.action!r}", status_code=3, grpc_status="INVALID_ARGUMENT" + ) + + body = json.loads(request.body) if request.body else {} + ticket_bytes = json.dumps(body, separators=(",", ":")).encode("utf-8") + ticket = flight.Ticket(ticket_bytes) + + try: + table = self._client.do_get(ticket, self._options()) + batches: list[pa.RecordBatch] = [] + while True: + try: + chunk, _metadata = table.read_chunk() + batches.append(chunk) + except StopIteration: + break + except flight.FlightError as ex: + raise TransportConnectionError(f"Error reading stream: {ex}") from ex + + return _to_bytes(pa.Table.from_batches(batches)) # validate schema + batches can be serialized + + except flight.FlightUnauthenticatedError as ex: + raise TransportAuthenticationError(str(ex)) from ex + except flight.FlightUnauthorizedError as ex: + raise TransportAuthorizationError(str(ex)) from ex + except flight.FlightUnavailableError as ex: + raise TransportConnectionError(str(ex)) from ex + except flight.FlightInternalError as ex: + raise TransportStatusError(str(ex), status_code=13, grpc_status="INTERNAL") from ex + except flight.FlightError as ex: + raise TransportConnectionError(str(ex)) from ex + except Exception as ex: + raise TransportConnectionError(f"Unexpected error during do_get: {ex}") from ex + def close(self) -> None: self._client.close() diff --git a/dataconnect/transport/base.py b/dataconnect/transport/base.py index df12c50..fe8b757 100644 --- a/dataconnect/transport/base.py +++ b/dataconnect/transport/base.py @@ -9,7 +9,7 @@ from abc import ABC, abstractmethod -from dataconnect.transport.models import ResourceInfo, ResourceQuery +from dataconnect.transport.models import ResourceInfo, ResourceQuery, DataTable class Transport(ABC): @@ -23,6 +23,14 @@ def list_resources(self, request: ResourceQuery) -> list[ResourceInfo]: service layer's responsibility. """ + @abstractmethod + def do_get(self, request: ResourceQuery) -> DataTable: + """Retrieve data for a single endpoint ticket. + + Reads all chunks from the server stream and returns a single + ``DataTable`` containing the complete result set. + """ + @abstractmethod def close(self) -> None: """Close the transport connection.""" diff --git a/dataconnect/transport/models.py b/dataconnect/transport/models.py index 65ea576..b7074b4 100644 --- a/dataconnect/transport/models.py +++ b/dataconnect/transport/models.py @@ -38,3 +38,13 @@ class ResourceInfo: endpoints: list[DataRef] total_records: int schema_bytes: bytes + +@dataclass(frozen=True) +class DataTable: + """Technology-agnostic representation of a fetched data result. + + ``schema_bytes`` holds the Arrow IPC-serialized schema. + ``ipc_bytes`` holds the full Arrow IPC stream (schema + all batches). + """ + schema_bytes: bytes + ipc_bytes: bytes \ No newline at end of file