From 3485ca5ac8fc69f624138613065673cbcc89e16a Mon Sep 17 00:00:00 2001 From: Tom Galindo <98626996+thommodin@users.noreply.github.com> Date: Fri, 3 Jul 2026 09:13:27 +1000 Subject: [PATCH 1/4] controllers done --- src/data_index/extract.py | 2 +- src/data_index/protocols.py | 9 ++++++--- src/data_index/transform.py | 12 ++++++------ 3 files changed, 13 insertions(+), 10 deletions(-) diff --git a/src/data_index/extract.py b/src/data_index/extract.py index 4d8185f..b45799f 100644 --- a/src/data_index/extract.py +++ b/src/data_index/extract.py @@ -32,7 +32,7 @@ def extract( # Return empty list if no object_references passed in if not object_references: logger.warning("extract called with no object references!") - return list() + return (list(), list()) # Count occurrences using the built-in versioned URI generator uri_counts = collections.Counter( diff --git a/src/data_index/protocols.py b/src/data_index/protocols.py index 3285968..906bf04 100644 --- a/src/data_index/protocols.py +++ b/src/data_index/protocols.py @@ -86,8 +86,11 @@ def to_compressed_base64_table( return compressed_base64_table - @staticmethod - def from_compressed_base64_table(base64_str: str) -> list[typing.Self]: + @classmethod + def from_compressed_base64_table( + cls, + base64_str: str, + ) -> list[typing.Self]: if not base64_str: return [] @@ -99,7 +102,7 @@ def from_compressed_base64_table(base64_str: str) -> list[typing.Self]: df = polars.read_ipc(buffer) # Reconstruct dataclass instances from the rows - return [ObjectReference(**row) for row in df.to_dicts()] + return [cls(**row) for row in df.to_dicts()] @dataclasses.dataclass( diff --git a/src/data_index/transform.py b/src/data_index/transform.py index 3412cb5..dd72fe4 100644 --- a/src/data_index/transform.py +++ b/src/data_index/transform.py @@ -12,7 +12,7 @@ def _transform_staged_object( staged_object: data_index.protocols.StagedObject, extractor: data_index.protocols.MetadataExtractor, - logger: logging.Logger, + logger: logging.Logger | logging.LoggerAdapter, ) -> data_index.protocols.ExtractedObject | data_index.protocols.DeadLetter: # Attempt to extract the metadata from the object @@ -38,9 +38,9 @@ def _transform_staged_object( def _transform_staged_objects( staged_objects: list[data_index.protocols.StagedObject], extractor: data_index.protocols.MetadataExtractor, - logger: logging.Logger, + logger: logging.Logger | logging.LoggerAdapter, ) -> tuple[ - list[data_index.protocols.ExtractedObject, list[data_index.protocols.DeadLetter]] + list[data_index.protocols.ExtractedObject], list[data_index.protocols.DeadLetter] ]: """ Populate all ObjectReferences with disk xarray handles. @@ -74,10 +74,10 @@ def _transform_staged_objects( def _concurrent_transform_staged_objects( staged_objects: list[data_index.protocols.StagedObject], extractor: data_index.protocols.MetadataExtractor, - logger: logging.Logger, + logger: logging.Logger | logging.LoggerAdapter, max_workers: int = 8, ) -> tuple[ - list[data_index.protocols.ExtractedObject, list[data_index.protocols.DeadLetter]] + list[data_index.protocols.ExtractedObject], list[data_index.protocols.DeadLetter] ]: """ Populate all ObjectReferences with disk xarray handles. @@ -142,7 +142,7 @@ def transform( # Return empty list if no object_references passed in if not staged_objects: logger.warning("transform called with no staged objects!") - return list() + return (list(), list()) logger.info("Running extraction sequentially...") From a2e075c94d5c76a72c70513928ad6bf7a418e7b9 Mon Sep 17 00:00:00 2001 From: Tom Galindo <98626996+thommodin@users.noreply.github.com> Date: Fri, 3 Jul 2026 09:21:03 +1000 Subject: [PATCH 2/4] fix protocols, fetchers --- src/data_index/file_fetcher/fsspec_fetcher.py | 2 +- src/data_index/file_fetcher/obstore_fetcher.py | 7 +++++-- src/data_index/protocols.py | 1 - 3 files changed, 6 insertions(+), 4 deletions(-) diff --git a/src/data_index/file_fetcher/fsspec_fetcher.py b/src/data_index/file_fetcher/fsspec_fetcher.py index d63c510..aedef09 100644 --- a/src/data_index/file_fetcher/fsspec_fetcher.py +++ b/src/data_index/file_fetcher/fsspec_fetcher.py @@ -39,7 +39,7 @@ def object_reference_to_staged_object( def fetch( self, object_references: list[data_index.protocols.ObjectReference] ) -> tuple[ - list[data_index.protocols.StagedObject, list[data_index.protocols.DeadLetter]] + list[data_index.protocols.StagedObject], list[data_index.protocols.DeadLetter] ]: staged_objects = [ diff --git a/src/data_index/file_fetcher/obstore_fetcher.py b/src/data_index/file_fetcher/obstore_fetcher.py index 75e0e59..ca63264 100644 --- a/src/data_index/file_fetcher/obstore_fetcher.py +++ b/src/data_index/file_fetcher/obstore_fetcher.py @@ -10,6 +10,9 @@ import data_index.protocols import data_index.xarray_handle +if typing.TYPE_CHECKING: + from obstore import GetOptions + class ObstoreFetcher(pydantic.BaseModel): type: typing.Literal["obstore_fetcher"] = pydantic.Field(default="obstore_fetcher") @@ -93,7 +96,7 @@ def get_stream( Convert an ObjectReference into a stream generator. """ - options = dict() + options: GetOptions = dict() # If version id if object_reference.version_id: @@ -107,7 +110,7 @@ def get_stream( def fetch( self, object_references: list[data_index.protocols.ObjectReference] ) -> tuple[ - list[data_index.protocols.StagedObject, list[data_index.protocols.DeadLetter]] + list[data_index.protocols.StagedObject], list[data_index.protocols.DeadLetter] ]: """ Populate all ObjectReferences with disk xarray handles. diff --git a/src/data_index/protocols.py b/src/data_index/protocols.py index 906bf04..0362b23 100644 --- a/src/data_index/protocols.py +++ b/src/data_index/protocols.py @@ -179,7 +179,6 @@ class ExtractionResult: @typing.runtime_checkable class XarrayHandle(typing.Protocol): - object_ref: ObjectReference file_format: str | None @property From e282438bebf1d77b89ed416d5c05af45ea6a8bd3 Mon Sep 17 00:00:00 2001 From: Tom Galindo <98626996+thommodin@users.noreply.github.com> Date: Fri, 3 Jul 2026 09:29:59 +1000 Subject: [PATCH 3/4] fix iceberg --- src/data_index/inventory_source/iceberg_table.py | 4 +--- 1 file changed, 1 insertion(+), 3 deletions(-) diff --git a/src/data_index/inventory_source/iceberg_table.py b/src/data_index/inventory_source/iceberg_table.py index 4941e2b..11b9f66 100644 --- a/src/data_index/inventory_source/iceberg_table.py +++ b/src/data_index/inventory_source/iceberg_table.py @@ -52,9 +52,7 @@ class IcebergTableFacilitySubsetInventorySource(IcebergTableInventorySource): subset_per_facility: int = pydantic.Field(default=10_000, ge=1) def inventory(self) -> polars.DataFrame: - df = self._scan( - selected_fields=("bucket", "key", "version_id", "size", "facility") - ) + df = self._scan() if df.is_empty(): return self._empty_inventory() From 41ac2c040e4631684fd5cea91d82fcd3e0200ac9 Mon Sep 17 00:00:00 2001 From: Tom Galindo <98626996+thommodin@users.noreply.github.com> Date: Fri, 3 Jul 2026 09:55:23 +1000 Subject: [PATCH 4/4] upgrade defaults --- src/data_index/runners/defaults.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/src/data_index/runners/defaults.py b/src/data_index/runners/defaults.py index d2a5240..1fe0ae4 100644 --- a/src/data_index/runners/defaults.py +++ b/src/data_index/runners/defaults.py @@ -49,7 +49,7 @@ table_config=_INVENTORY_TABLE_CONFIG, table_scan_config=IcebergTableScanConfig( row_filter="key LIKE 'IMOS/SOOP/%' OR key LIKE 'IMOS/AATAMS/%' OR key LIKE 'IMOS/ANMN/%' OR key LIKE 'IMOS/FAIMMS/%' OR key LIKE 'IMOS/OceanCurrent/%' OR key LIKE 'IMOS/DWM/%' OR key LIKE 'IMOS/AUV/%' OR key LIKE 'IMOS/COASTAL-WAVE-BUOYS/%' OR key LIKE 'IMOS/NTP/%' OR key LIKE 'IMOS/ANFOG/%' OR key LIKE 'IMOS/eMII/%'", - selected_fields=["bucket", "key", "version_id", "size"], + selected_fields=("bucket", "key", "version_id", "size"), ), )