Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion pyproject.toml
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ build-backend = "hatchling.build"

[project]
name = "dri-timeseries-processor"
version = "0.7.31"
version = "0.7.33"
description = "Timeseries processor service."
readme = "README.md"
license = { file = "LICENSE" }
Expand Down
22 changes: 17 additions & 5 deletions src/dritimeseriesprocessor/dag/dataset_dependency_graph.py
Original file line number Diff line number Diff line change
Expand Up @@ -301,7 +301,7 @@ def _fetch_configs_for_dataset(self, dataset_ids: list[str]) -> dict[str, list[D
return {}

response = self.metadata_router.fetch_processing_configs(dataset_ids)
dataset_configs = self._build_processing_configs(response)
dataset_configs = self._build_processing_configs(response, dataset_ids)
return dataset_configs

def _fetch_root_datasets_by_ids(self, dataset_ids: list[str]) -> list[TimeSeriesContainer]:
Expand Down Expand Up @@ -418,21 +418,33 @@ def _build_dataset_containers(self, dataset_response: TimeSeriesDatasetResponse)
return all_containers

def _build_processing_configs(
self, dataset_response: DataProcessingConfiguration
self, dataset_response: DataProcessingConfiguration, requested_ids: list[str]
) -> dict[str, list[DataProcessingConfig]]:
"""Parse an API response containing processing configuration items and return them grouped by timeseries ID.

A configuration can apply to several datasets. Each of those datasets gets its own copy of the config, so the
grouping does not depend on the order the API happens to list them in.

Args:
dataset_response: The DataProcessingConfig representing a set of processing configs.
requested_ids: The dataset IDs the configs were requested for.

Returns:
A dictionary keyed by timeseries ID, with value as the list of associated processing configs.
"""
requested = set(requested_ids)
dataset_configs = defaultdict(list)
for item in dataset_response.items:
mapped_config = map_processing_config_item(item, self.site_metadata)
self._get_deployment_attributes(mapped_config)
dataset_configs[mapped_config.ts_id].append(mapped_config)
applies_to_requested = [applies_to for applies_to in item.applies_to_dataset if applies_to.id in requested]
if not applies_to_requested:
logger.warning(
f"Config {item.id} applies to none of the requested datasets: "
f"{[applies_to.id for applies_to in item.applies_to_dataset]}"
)
for applies_to in applies_to_requested:
mapped_config = map_processing_config_item(item, applies_to, self.site_metadata)
self._get_deployment_attributes(mapped_config)
dataset_configs[mapped_config.ts_id].append(mapped_config)
return dataset_configs

def _get_deployment_attributes(self, mapped_config: DataProcessingConfig) -> None:
Expand Down
13 changes: 10 additions & 3 deletions src/dritimeseriesprocessor/models/mappers/api_to_domain.py
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@

from dritimeseriesprocessor.models.api_models.annotation import HasAnnotationItem
from dritimeseriesprocessor.models.api_models.data_processing_configuration import (
AppliesToDataset,
DataProcessingConfigurationItem,
)
from dritimeseriesprocessor.models.api_models.dataset_observation import ObservationDatasetItem
Expand Down Expand Up @@ -89,23 +90,29 @@ def map_dataset_item(


def map_processing_config_item(
item: DataProcessingConfigurationItem, all_site_metadata: dict[str, SiteMetadata]
item: DataProcessingConfigurationItem,
applies_to: AppliesToDataset,
all_site_metadata: dict[str, SiteMetadata],
) -> DataProcessingConfig:
"""Map a DataProcessingConfigurationItem to a ProcessingConfig domain model.

Some processing configurations will have "site_attribute" parameters that require fetching this metadata
key from the site metadata.

A configuration can apply to more than one dataset. We use `applies_to` to pick which of those datasets this
config is being mapped for.

Args:
item: The validated DataProcessingConfigurationItem from the API.
applies_to: The dataset, from the item's `appliesToDataset` list, to map this config for.
all_site_metadata: Metadata for sites.

Returns:
A ProcessingConfig domain object containing annotations and a list of MethodConfig objects which provide
specific method configurations for use in the processing pipeline
"""
ts_id = item.applies_to_dataset[0].id
site_id = item.applies_to_dataset[0].originating_site[0].id
ts_id = applies_to.id
site_id = applies_to.originating_site[0].id
config_type = ConfigurationType(extract_uri_id(item.type.id))
annotations = extract_annotations(item.has_annotation)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -114,6 +114,7 @@ def create_items_list(ts_ids: str | list) -> list:
item = MagicMock()
item.__getitem__ = MagicMock(side_effect=lambda key, _id=ts_id: _id)
item.originating_site = [MagicMock(id=ts_id)]
item.applies_to_dataset = [MagicMock(id=ts_id)]
items.append(item)
return items

Expand Down Expand Up @@ -175,7 +176,7 @@ def monkeypatch_mappers(monkeypatch: pytest.MonkeyPatch) -> None:

monkeypatch.setattr(
"dritimeseriesprocessor.dag.dataset_dependency_graph.map_processing_config_item",
lambda item, _: make_processing_config_container(item["@id"]),
lambda _item, applies_to, _site_metadata: make_processing_config_container(applies_to.id),
)

monkeypatch.setattr(
Expand Down Expand Up @@ -221,6 +222,35 @@ def test_fetch_configs_for_single_dataset(self, monkeypatch: pytest.MonkeyPatch)
assert result == {ts_id: [container]}
assert mock_router.fetch_processing_configs.call_count == 1

@pytest.mark.parametrize("applies_to_ids", [["ds1", "ds2"], ["ds2", "ds1"]])
def test_fetch_configs_for_shared_config(self, applies_to_ids: list[str], monkeypatch: pytest.MonkeyPatch) -> None:
"""Tests that a config applying to several datasets is returned for each of them, whatever order the
metadata API lists those datasets in."""
mock_router = setup_mocks(["ds1", "ds2"], monkeypatch)
shared_item = create_items_list("shared_config")[0]
shared_item.applies_to_dataset = [MagicMock(id=ts_id) for ts_id in applies_to_ids]
mock_router.fetch_processing_configs.return_value = SimpleNamespace(items=[shared_item])
builder = DatasetDependencyGraph(mock_router, MagicMock(), MagicMock(), MagicMock())

result = builder._fetch_configs_for_dataset(["ds1", "ds2"])

assert result == {
"ds1": [make_processing_config_container("ds1")],
"ds2": [make_processing_config_container("ds2")],
}

def test_fetch_configs_ignores_unrequested_datasets(self, monkeypatch: pytest.MonkeyPatch) -> None:
"""Tests that a config is not returned for datasets that were not asked for."""
mock_router = setup_mocks(["ds1", "ds2"], monkeypatch)
shared_item = create_items_list("shared_config")[0]
shared_item.applies_to_dataset = [MagicMock(id="ds2")]
mock_router.fetch_processing_configs.return_value = SimpleNamespace(items=[shared_item])
builder = DatasetDependencyGraph(mock_router, MagicMock(), MagicMock(), MagicMock())

result = builder._fetch_configs_for_dataset(["ds1"])

assert result == {}

@pytest.mark.parametrize(
"sites",
[
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -625,10 +625,12 @@ def test_value_series_annotation_without_direct_values(self) -> None:

class TestMapProcessingConfigItem:
def test_qc_processing_config(self) -> None:
"""Tests that a QC configuration item maps to a DataProcessingConfig with its method and parameters."""
filename = TEST_DATA_API_VALID / "data_processing_configuration" / "cosmos_bunny_swin_30min_qc.json"
api_model = valid_parses(load_json_file, filename, DataProcessingConfiguration)
item = api_model.items[0]

result = map_processing_config_item(api_model.items[0], MagicMock())
result = map_processing_config_item(item, item.applies_to_dataset[0], MagicMock())

expected = DataProcessingConfig(
ts_id="http://fdri.ceh.ac.uk/id/dataset/cosmos-bunny-swin_30min_processed",
Expand All @@ -647,10 +649,12 @@ def test_qc_processing_config(self) -> None:
assert result == expected

def test_infill_processing_config(self) -> None:
"""Tests that an infill configuration item maps to a DataProcessingConfig with its method and parameters."""
filename = TEST_DATA_API_VALID / "data_processing_configuration" / "cosmos_bunny_swin_30min_infill.json"
api_model = valid_parses(load_json_file, filename, DataProcessingConfiguration)
item = api_model.items[0]

result = map_processing_config_item(api_model.items[0], MagicMock())
result = map_processing_config_item(item, item.applies_to_dataset[0], MagicMock())

expected = DataProcessingConfig(
ts_id="http://fdri.ceh.ac.uk/id/dataset/cosmos-bunny-swin_30min_processed",
Expand All @@ -669,10 +673,12 @@ def test_infill_processing_config(self) -> None:
assert result == expected

def test_correction_processing_config(self) -> None:
"""Tests that a correction configuration item maps to a DataProcessingConfig with its observation interval."""
filename = TEST_DATA_API_VALID / "data_processing_configuration" / "cosmos_bunny_swin_30min_correction.json"
api_model = valid_parses(load_json_file, filename, DataProcessingConfiguration)
item = api_model.items[0]

result = map_processing_config_item(api_model.items[0], MagicMock())
result = map_processing_config_item(item, item.applies_to_dataset[0], MagicMock())

expected = DataProcessingConfig(
ts_id="http://fdri.ceh.ac.uk/id/dataset/cosmos-bunny-swin_30min_processed",
Expand Down
Loading
Loading