diff --git a/docs/api.rst b/docs/api.rst index d9e7223..32b07e4 100644 --- a/docs/api.rst +++ b/docs/api.rst @@ -15,6 +15,7 @@ API Reference fetchers/czech fetchers/france fetchers/germany_berlin + fetchers/greece fetchers/japan fetchers/lithuania fetchers/norway @@ -25,4 +26,4 @@ API Reference fetchers/spain fetchers/uk_ea fetchers/uk_nrfa - fetchers/usa \ No newline at end of file + fetchers/usa diff --git a/docs/fetchers/greece.rst b/docs/fetchers/greece.rst new file mode 100644 index 0000000..306e2ab --- /dev/null +++ b/docs/fetchers/greece.rst @@ -0,0 +1,5 @@ +Greece Fetcher +============== + +.. automodule:: rivretrieve.greece + :members: diff --git a/examples/test_greece_fetcher.py b/examples/test_greece_fetcher.py new file mode 100644 index 0000000..77d38d3 --- /dev/null +++ b/examples/test_greece_fetcher.py @@ -0,0 +1,36 @@ +import logging + +import matplotlib.pyplot as plt + +from rivretrieve import GreeceFetcher, constants + +logging.basicConfig(level=logging.INFO) + +gauge_id = "1458" +variables = [ + constants.STAGE_DAILY_MEAN, + constants.DISCHARGE_DAILY_MEAN, +] +start_date = "2025-01-01" +end_date = "2025-01-07" + +fetcher = GreeceFetcher() + +for variable in variables: + data = fetcher.get_data(gauge_id=gauge_id, variable=variable, start_date=start_date, end_date=end_date) + if data.empty: + print(f"No data found for {gauge_id} ({variable})") + continue + + print(data.head()) + plt.figure(figsize=(12, 6)) + plt.plot(data.index, data[variable], label=f"{gauge_id} - {variable}") + plt.xlabel(constants.TIME_INDEX) + plt.ylabel(variable) + plt.title(f"Greece OpenHI River Data ({gauge_id})") + plt.legend() + plt.grid(True) + plt.tight_layout() + plot_path = f"greece_{variable}_plot.png" + plt.savefig(plot_path) + print(f"Plot saved to {plot_path}") diff --git a/rivretrieve/__init__.py b/rivretrieve/__init__.py index ae2e150..59ddcc3 100644 --- a/rivretrieve/__init__.py +++ b/rivretrieve/__init__.py @@ -8,6 +8,7 @@ from .czech import CzechFetcher from .france import FranceFetcher from .germany_berlin import GermanyBerlinFetcher +from .greece import GreeceFetcher from .japan import JapanFetcher from .lithuania import LithuaniaFetcher from .norway import NorwayFetcher diff --git a/rivretrieve/cached_site_data/greece_sites.csv b/rivretrieve/cached_site_data/greece_sites.csv new file mode 100644 index 0000000..01b2b32 --- /dev/null +++ b/rivretrieve/cached_site_data/greece_sites.csv @@ -0,0 +1,44 @@ +gauge_id,station_name,river,latitude,longitude,altitude,area,country,source,station_code,display_timezone,owner_id,owner_name,provider_start_date,provider_end_date +1356,Καρβελιώτης - Karveliotis,,37.073479,22.223609,598.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,2011-12-16, +1458,Ανθήλη,,38.856109,22.466853,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +1461,Βαρυμπόμπη,,38.107278,23.807221,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,14,Εθνικό Μετσόβιο Πολυτεχνείο - Σχολή Μηχανικών Μεταλλείων - Μεταλλουργών,2018-06-08, +1462,Δεκέλεια,,38.092422,23.770555,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,14,Εθνικό Μετσόβιο Πολυτεχνείο - Σχολή Μηχανικών Μεταλλείων - Μεταλλουργών,2018-06-08, +1463,Μοναστήρι,,38.059284,23.746475,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,14,Εθνικό Μετσόβιο Πολυτεχνείο - Σχολή Μηχανικών Μεταλλείων - Μεταλλουργών,2018-06-08, +1464,Κόκκινος Μύλος,,38.045492,23.73706,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,14,Εθνικό Μετσόβιο Πολυτεχνείο - Σχολή Μηχανικών Μεταλλείων - Μεταλλουργών,2018-06-08, +1466,Ρέντης,,37.960928,23.676186,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,14,Εθνικό Μετσόβιο Πολυτεχνείο - Σχολή Μηχανικών Μεταλλείων - Μεταλλουργών,2018-06-08, +1481,Σαρανταπόταμος (Γύρα Στεφάνης) - Sarantapotamos,,38.132831,23.533028,157.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,2018-07-05, +1482,Νέδοντας Καλαμάτα (Γέφυρα Μεγάρου Χορού) - Nedon Kalamata,,37.038592,22.1075,75.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,2018-07-05, +1484,Arta's Bridge,,39.15104,20.975265,,,Greece,OpenHI / Enhydris public API,Arta,Etc/GMT-2,21,Laboratory of Knowledge and Intelligent Computing (KIC),, +1485,Neoxori Bridge,,39.070688,21.025465,,,Greece,OpenHI / Enhydris public API,Neoxori,Etc/GMT-2,21,Laboratory of Knowledge and Intelligent Computing (KIC),, +1486,Αλαγονία (Νερόμυλος Ρεντίφη) - Alagonia,,37.104089,22.2325,562.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,2018-07-05, +1487,Νέδουσα - Nedousa,,37.116714,22.203611,392.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,2018-07-05, +1488,Σέλας (Γέφυρα Γκολφ Costa Navarino) - Selas,,36.99445,21.658197,9.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,, +1534,Αλαμάνα,,38.8125,22.4952,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +1544,Νομή,,39.52657,21.93833,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +1545,Γ. Γιάννουλη,,39.65246,22.4078,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +1546,Τέμπη,,39.89675,22.6152,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +2018,Επιτάλιο,,37.64256,21.47648,1.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +2019,Άσπρα Σπίτια,,37.586411,21.790867,56.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +2020,Μεσοχώρα Κατάντη,,39.420104,21.262608,581.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +2022,Ευρώτας - Eurotas,,37.130856,22.399872,224.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2039,Λούσιος (Γέφυρα Ατσίχολου) - Loussios (Atsicholos Bridge),,37.511414,22.038225,230.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2075,Σαρανταπόταμος (Οινόη) - Oinoe,,38.165317,23.399436,333.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2191,Μαυροζούμαινα - Mavrozoumena,,37.141406,21.988792,20.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2192,Χαλάνδρι - Chalandri,,38.024,23.795431,167.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2193,Φιλοθέη - Filothei,,38.021644,23.785278,161.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2194,Ποδονίφτης - Podoniftis,,38.025381,23.730247,74.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2195,Προφήτης Δανιήλ - Profitis Daniel,,37.975725,23.690911,21.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2196,Κελεφίνα - Kelefina,,37.117986,22.453142,251.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +2208,Τ2 Τ21 Α Ζώνη Λούρου (ΓΟΕΒ Π Άρτας),,39.07878,20.88414,-1.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,24,Περιφέρεια Ηπείρου,, +2209,Τ1 Βίγλας (ΓΟΕΒ Π Άρτας),,39.08824,20.87594,-1.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,24,Περιφέρεια Ηπείρου,, +2210,Τ2 Α Ζώνη Λούρου Τ0 Σαλαώρας (ΓΟΕΒ Π Άρτας),,39.07907,20.88588,-1.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,24,Περιφέρεια Ηπείρου,, +2213,Πάμισος - Pamissos,,37.051972,22.020492,5.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,17,Εθνικό Αστεροσκοπείο Αθηνών,, +28071,Γούρια,,38.480525,21.260677,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +28081,Μεσοχώρα Κατάντη (new),,39.384321,21.272844,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +28082,Γούρια (παροχή),,38.48039,21.26099,,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +28083,Μυλοπόταμος - Mylopotamos,,36.24548,22.94454,255.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +28084,Καραβάς - Karavas,,36.34579,22.94914,108.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT,17,Εθνικό Αστεροσκοπείο Αθηνών,, +4192,Τριχωνίδα,,38.590884,21.600644,16.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,11,Ελληνικό Κέντρο Θαλάσσιων Ερευνών,, +8424,Μάνδρα Ρέμα,,38.098118,23.456293,215.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,25,"Εθνικό Αστεροσκοπείο Αθηνών, ΙΑΑΔΕΤ, Κέντρο Επιστημών Παρατήρησης της Γης και Δορυφορικής Τηλεπισκόπησης BEYOND",, +8425,Μάνδρα Εκτροπή,,38.080206,23.48023,118.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,25,"Εθνικό Αστεροσκοπείο Αθηνών, ΙΑΑΔΕΤ, Κέντρο Επιστημών Παρατήρησης της Γης και Δορυφορικής Τηλεπισκόπησης BEYOND",, +8426,Μάνδρα Κόμβος,,38.072345,23.509526,61.0,,Greece,OpenHI / Enhydris public API,,Etc/GMT-2,25,"Εθνικό Αστεροσκοπείο Αθηνών, ΙΑΑΔΕΤ, Κέντρο Επιστημών Παρατήρησης της Γης και Δορυφορικής Τηλεπισκόπησης BEYOND",, diff --git a/rivretrieve/greece.py b/rivretrieve/greece.py new file mode 100644 index 0000000..52840ff --- /dev/null +++ b/rivretrieve/greece.py @@ -0,0 +1,598 @@ +"""Fetcher for Greece river gauge data from the OpenHI / Enhydris public API.""" + +from __future__ import annotations + +import logging +import re +from io import StringIO +from typing import Any, Optional + +import numpy as np +import pandas as pd +import requests + +from . import base, constants, utils + +logger = logging.getLogger(__name__) + + +class GreeceFetcher(base.RiverDataFetcher): + """Fetches Greek river gauge data from the public OpenHI / Enhydris API. + + Data source: + https://openhi.net/en/ + + Supported variables: + - 'discharge_daily_mean' (m³/s) + - 'discharge_instantaneous' (m³/s) + - 'stage_daily_mean' (m) + - 'stage_instantaneous' (m) + + Data description and API: + - see https://system.openhi.net/api/ + - see https://enhydris.readthedocs.io/en/latest/dev/webservice-api.html + + Terms of use: + - see https://openhi.net/licence/ + """ + + BASE_URL = "https://system.openhi.net/api" + SOURCE = "OpenHI / Enhydris public API" + COUNTRY = "Greece" + STATION_SEARCH_QUERIES = ( + "ts_only: variable:stage", + "ts_only: variable:level", + "ts_only: variable:discharge", + ) + SENTINEL_VALUES = {-999.0, -9999.0, -6999.0} + TYPE_PRIORITY = { + "Checked": 0, + "Initial": 1, + "Aggregated": 2, + } + VARIABLE_MAP = { + constants.DISCHARGE_DAILY_MEAN: { + "kind": constants.DISCHARGE, + "aggregate_daily": True, + "provider_variable_ids": {2}, + }, + constants.DISCHARGE_INSTANT: { + "kind": constants.DISCHARGE, + "aggregate_daily": False, + "provider_variable_ids": {2}, + }, + constants.STAGE_DAILY_MEAN: { + "kind": constants.STAGE, + "aggregate_daily": True, + "provider_variable_ids": {14, 5710}, + }, + constants.STAGE_INSTANT: { + "kind": constants.STAGE, + "aggregate_daily": False, + "provider_variable_ids": {14, 5710}, + }, + } + + def __init__(self): + self._session = utils.requests_retry_session( + retries=6, + backoff_factor=1, + status_forcelist=(429, 500, 502, 503, 504), + ) + self._station_query_cache: dict[str, list[dict[str, Any]]] = {} + self._timeseries_group_cache: dict[str, list[dict[str, Any]]] = {} + self._timeseries_cache: dict[tuple[str, str], list[dict[str, Any]]] = {} + self._variable_cache: dict[int, str] = {} + self._unit_cache: dict[int, str] = {} + self._organization_cache: dict[int, Optional[str]] = {} + + @staticmethod + def get_cached_metadata() -> pd.DataFrame: + """Retrieves cached Greece gauge metadata.""" + return utils.load_cached_metadata_csv("greece") + + @staticmethod + def get_available_variables() -> tuple[str, ...]: + return tuple(GreeceFetcher.VARIABLE_MAP.keys()) + + @staticmethod + def _empty_data_frame(variable: str) -> pd.DataFrame: + return pd.DataFrame(columns=[constants.TIME_INDEX, variable]).set_index(constants.TIME_INDEX) + + @staticmethod + def _empty_metadata_frame() -> pd.DataFrame: + columns = [ + constants.GAUGE_ID, + constants.STATION_NAME, + constants.RIVER, + constants.LATITUDE, + constants.LONGITUDE, + constants.ALTITUDE, + constants.AREA, + constants.COUNTRY, + constants.SOURCE, + "station_code", + "display_timezone", + "owner_id", + "owner_name", + "provider_start_date", + "provider_end_date", + ] + return pd.DataFrame(columns=columns).set_index(constants.GAUGE_ID) + + @classmethod + def _absolute_url(cls, path_or_url: str) -> str: + if path_or_url.startswith("http://") or path_or_url.startswith("https://"): + return path_or_url + return f"{cls.BASE_URL}/{path_or_url.lstrip('/')}" + + def _request_json(self, path_or_url: str, params: Optional[dict[str, Any]] = None) -> Any: + url = self._absolute_url(path_or_url) + response = self._session.get(url, params=params, timeout=60) + response.raise_for_status() + return response.json() + + def _request_text(self, path_or_url: str, params: Optional[dict[str, Any]] = None) -> str: + url = self._absolute_url(path_or_url) + response = self._session.get(url, params=params, timeout=60) + response.raise_for_status() + return response.text + + def _fetch_paginated(self, path_or_url: str, params: Optional[dict[str, Any]] = None) -> list[dict[str, Any]]: + url = self._absolute_url(path_or_url) + current_params = params.copy() if params else None + results: list[dict[str, Any]] = [] + + while url: + payload = self._request_json(url, params=current_params) + if isinstance(payload, dict): + page_results = payload.get("results", []) + if isinstance(page_results, list): + results.extend([record for record in page_results if isinstance(record, dict)]) + url = payload.get("next") + current_params = None + continue + break + + return results + + @staticmethod + def _parse_geom(geom: Any) -> tuple[float, float]: + if geom is None: + return np.nan, np.nan + + match = re.search(r"POINT\s*\(([-+]?\d+(?:\.\d+)?)\s+([-+]?\d+(?:\.\d+)?)\)", str(geom)) + if not match: + return np.nan, np.nan + + longitude = pd.to_numeric(match.group(1), errors="coerce") + latitude = pd.to_numeric(match.group(2), errors="coerce") + return float(latitude), float(longitude) + + @staticmethod + def _normalized_text(value: Any) -> str: + return str(value or "").strip().casefold() + + @staticmethod + def _timestamp_sort_value(value: Any) -> int: + parsed = pd.to_datetime(value, errors="coerce") + if pd.isna(parsed): + return 0 + return -int(parsed.value) + + def _search_stations(self, query: str) -> list[dict[str, Any]]: + if query not in self._station_query_cache: + self._station_query_cache[query] = self._fetch_paginated("/stations/", params={"q": query}) + return [record.copy() for record in self._station_query_cache[query]] + + def _get_timeseries_groups(self, station_id: str) -> list[dict[str, Any]]: + station_id = str(station_id).strip() + if station_id not in self._timeseries_group_cache: + self._timeseries_group_cache[station_id] = self._fetch_paginated( + f"/stations/{station_id}/timeseriesgroups/" + ) + return [record.copy() for record in self._timeseries_group_cache[station_id]] + + def _get_timeseries(self, station_id: str, group_id: str) -> list[dict[str, Any]]: + key = (str(station_id).strip(), str(group_id).strip()) + if key not in self._timeseries_cache: + self._timeseries_cache[key] = self._fetch_paginated( + f"/stations/{key[0]}/timeseriesgroups/{key[1]}/timeseries/" + ) + return [record.copy() for record in self._timeseries_cache[key]] + + def _get_variable_name(self, variable_id: Any) -> str: + numeric_id = pd.to_numeric(variable_id, errors="coerce") + if pd.isna(numeric_id): + return "" + + cache_key = int(numeric_id) + if cache_key not in self._variable_cache: + payload = self._request_json(f"/variables/{cache_key}/") + translations = payload.get("translations", {}) if isinstance(payload, dict) else {} + english = translations.get("en", {}) if isinstance(translations, dict) else {} + greek = translations.get("el", {}) if isinstance(translations, dict) else {} + description = english.get("descr") or greek.get("descr") or "" + self._variable_cache[cache_key] = str(description).strip() + + return self._variable_cache[cache_key] + + def _get_unit_symbol(self, unit_id: Any) -> str: + numeric_id = pd.to_numeric(unit_id, errors="coerce") + if pd.isna(numeric_id): + return "" + + cache_key = int(numeric_id) + if cache_key not in self._unit_cache: + payload = self._request_json(f"/units/{cache_key}/") + symbol = payload.get("symbol", "") if isinstance(payload, dict) else "" + self._unit_cache[cache_key] = str(symbol).strip() + + return self._unit_cache[cache_key] + + def _get_owner_name(self, owner_id: Any) -> Optional[str]: + numeric_id = pd.to_numeric(owner_id, errors="coerce") + if pd.isna(numeric_id): + return None + + cache_key = int(numeric_id) + if cache_key not in self._organization_cache: + payload = self._request_json(f"/organizations/{cache_key}/") + if not isinstance(payload, dict): + self._organization_cache[cache_key] = None + else: + name = payload.get("name") or payload.get("ordering_string") or payload.get("acronym") + self._organization_cache[cache_key] = str(name).strip() or None + + return self._organization_cache[cache_key] + + def get_metadata(self) -> pd.DataFrame: + """Fetches live metadata for stations with supported Greek variables. + + Merges the public OpenHI station searches for stage, water level, and + discharge stations and returns standardized metadata indexed by + ``constants.GAUGE_ID``. + """ + stations_by_id: dict[str, dict[str, Any]] = {} + + for query in self.STATION_SEARCH_QUERIES: + for station in self._search_stations(query): + station_id = str(station.get("id", "")).strip() + if station_id: + stations_by_id[station_id] = station + + if not stations_by_id: + return self._empty_metadata_frame() + + rows = [] + for gauge_id, station in stations_by_id.items(): + latitude, longitude = self._parse_geom(station.get("geom")) + owner_id = pd.to_numeric(station.get("owner"), errors="coerce") + rows.append( + { + constants.GAUGE_ID: gauge_id, + constants.STATION_NAME: station.get("name"), + constants.RIVER: np.nan, + constants.LATITUDE: latitude, + constants.LONGITUDE: longitude, + constants.ALTITUDE: pd.to_numeric(station.get("altitude"), errors="coerce"), + constants.AREA: np.nan, + constants.COUNTRY: self.COUNTRY, + constants.SOURCE: self.SOURCE, + "station_code": station.get("code") or None, + "display_timezone": station.get("display_timezone") or None, + "owner_id": owner_id, + "owner_name": self._get_owner_name(owner_id), + "provider_start_date": station.get("start_date") or None, + "provider_end_date": station.get("end_date") or None, + } + ) + + df = pd.DataFrame(rows) + if df.empty: + return self._empty_metadata_frame() + + df[constants.GAUGE_ID] = df[constants.GAUGE_ID].astype(str).str.strip() + df = df.drop_duplicates(subset=[constants.GAUGE_ID]).sort_values(constants.GAUGE_ID) + return df.set_index(constants.GAUGE_ID) + + def _group_kind(self, group: dict[str, Any]) -> Optional[str]: + variable_id = pd.to_numeric(group.get("variable"), errors="coerce") + if pd.notna(variable_id): + variable_id = int(variable_id) + if variable_id == 2: + return constants.DISCHARGE + if variable_id in {14, 5710}: + return constants.STAGE + + variable_name = self._normalized_text(self._get_variable_name(group.get("variable"))) + if "discharge" in variable_name: + return constants.DISCHARGE + if variable_name in {"stage", "water level"}: + return constants.STAGE + return None + + def _group_semantic_score(self, group_name: Any, kind: str, aggregate_daily: bool) -> int: + text = self._normalized_text(group_name) + base_tokens = {"discharge"} if kind == constants.DISCHARGE else {"stage", "water level", "level"} + has_base_name = not text or any(token in text for token in base_tokens) + has_mean = any(token in text for token in ("mean", "avg", "average", "daily")) + has_extreme = any(token in text for token in ("max", "maximum", "min", "minimum")) + has_old = "old" in text + + if has_extreme: + score = 4 + elif aggregate_daily and has_mean: + score = 0 + elif has_base_name: + score = 0 + elif not aggregate_daily and has_mean: + score = 2 + else: + score = 1 + + if has_old: + score += 2 + return score + + def _ranked_groups(self, gauge_id: str, variable: str) -> list[dict[str, Any]]: + config = self.VARIABLE_MAP[variable] + candidates = [ + group + for group in self._get_timeseries_groups(gauge_id) + if self._group_kind(group) == config["kind"] + and pd.to_numeric(group.get("variable"), errors="coerce") in config["provider_variable_ids"] + ] + return sorted( + candidates, + key=lambda group: ( + self._group_semantic_score(group.get("name"), config["kind"], config["aggregate_daily"]), + self._timestamp_sort_value(group.get("last_modified")), + ), + ) + + @staticmethod + def _time_step_minutes(time_step: Any) -> Optional[int]: + text = str(time_step or "").strip().lower() + if not text: + return None + + match = re.fullmatch(r"(\d+)\s*([a-z]+)", text) + if not match: + return None + + amount = int(match.group(1)) + unit = match.group(2) + multipliers = { + "m": 1, + "min": 1, + "mins": 1, + "minute": 1, + "minutes": 1, + "h": 60, + "hr": 60, + "hrs": 60, + "hour": 60, + "hours": 60, + "d": 1440, + "day": 1440, + "days": 1440, + } + if unit not in multipliers: + return None + return amount * multipliers[unit] + + @classmethod + def _timeseries_priority(cls, timeseries: dict[str, Any], aggregate_daily: bool) -> tuple[int, int, int]: + timeseries_type = str(timeseries.get("type") or "").strip() + type_rank = cls.TYPE_PRIORITY.get(timeseries_type, 99) + step_minutes = cls._time_step_minutes(timeseries.get("time_step")) + step_rank = step_minutes if step_minutes is not None else -1 + is_daily = step_minutes is not None and step_minutes >= 1440 + name = cls._normalized_text(timeseries.get("name")) + + if aggregate_daily: + if timeseries_type == "Aggregated" and is_daily and "mean" in name: + class_rank = 0 + elif timeseries_type == "Aggregated" and is_daily: + class_rank = 1 + elif timeseries_type == "Checked" and not is_daily: + class_rank = 2 + elif timeseries_type == "Initial" and not is_daily: + class_rank = 3 + elif timeseries_type == "Aggregated" and not is_daily: + class_rank = 4 + else: + class_rank = 5 + else: + if timeseries_type == "Checked" and not is_daily: + class_rank = 0 + elif timeseries_type == "Initial" and not is_daily: + class_rank = 1 + elif timeseries_type == "Aggregated" and not is_daily: + class_rank = 2 + elif timeseries_type == "Checked" and is_daily: + class_rank = 3 + elif timeseries_type == "Initial" and is_daily: + class_rank = 4 + else: + class_rank = 5 + + return ( + class_rank, + type_rank, + step_rank, + ) + + def _ranked_timeseries(self, gauge_id: str, group_id: str, aggregate_daily: bool) -> list[dict[str, Any]]: + series = [ + record + for record in self._get_timeseries(gauge_id, group_id) + if bool(record.get("publicly_available", True)) + ] + if not series: + return [] + + return sorted( + series, + key=lambda record: ( + self._timeseries_priority(record, aggregate_daily), + self._timestamp_sort_value(record.get("last_modified")), + ), + ) + + @staticmethod + def _unit_factor(kind: str, unit_symbol: str) -> float: + normalized_unit = str(unit_symbol or "").strip() + + if kind == constants.STAGE and normalized_unit == "cm": + return 0.01 + if kind == constants.DISCHARGE and normalized_unit in {"l/s", "L/s"}: + return 0.001 + return 1.0 + + @classmethod + def _parse_csv_payload(cls, payload: str, factor: float) -> pd.DataFrame: + if not payload.strip(): + return pd.DataFrame(columns=[constants.TIME_INDEX, "value"]) + + raw_df = pd.read_csv(StringIO(payload), header=None) + if raw_df.empty or raw_df.shape[1] < 2: + return pd.DataFrame(columns=[constants.TIME_INDEX, "value"]) + + result = pd.DataFrame( + { + constants.TIME_INDEX: pd.to_datetime(raw_df.iloc[:, 0], errors="coerce"), + "value": pd.to_numeric(raw_df.iloc[:, 1], errors="coerce"), + } + ).dropna(subset=[constants.TIME_INDEX, "value"]) + + if result.empty: + return result + + result.loc[result["value"].isin(cls.SENTINEL_VALUES), "value"] = np.nan + result = result.dropna(subset=["value"]) + if result.empty: + return result + + result["value"] = result["value"] * factor + return result + + def _download_data( + self, + gauge_id: str, + variable: str, + start_date: str, + end_date: str, + ) -> list[dict[str, Any]]: + config = self.VARIABLE_MAP[variable] + ranked_groups = self._ranked_groups(gauge_id, variable) + payloads: list[dict[str, Any]] = [] + + for group_rank, group in enumerate(ranked_groups): + group_id = str(group.get("id", "")).strip() + if not group_id: + continue + + ranked_series = self._ranked_timeseries(gauge_id, group_id, config["aggregate_daily"]) + if not ranked_series: + continue + + for series_rank, timeseries in enumerate(ranked_series): + timeseries_id = str(timeseries.get("id", "")).strip() + if not timeseries_id: + continue + + try: + payload = self._request_text( + f"/stations/{gauge_id}/timeseriesgroups/{group_id}/timeseries/{timeseries_id}/data/", + params={ + "fmt": "csv", + "start_date": start_date, + "end_date": end_date, + }, + ) + except requests.exceptions.RequestException as exc: + logger.warning( + f"Failed to fetch Greece data for station {gauge_id}, group {group_id}, " + f"timeseries {timeseries_id}: {exc}" + ) + continue + + if not payload.strip(): + continue + + payloads.append( + { + "payload": payload, + "group_rank": group_rank, + "series_rank": series_rank, + "unit_symbol": self._get_unit_symbol(group.get("unit_of_measurement")), + } + ) + + return payloads + + def _parse_data(self, gauge_id: str, raw_data: list[dict[str, Any]], variable: str) -> pd.DataFrame: + config = self.VARIABLE_MAP[variable] + if not raw_data: + return self._empty_data_frame(variable) + + frames = [] + for payload in raw_data: + parsed = self._parse_csv_payload( + payload.get("payload", ""), + factor=self._unit_factor(config["kind"], payload.get("unit_symbol", "")), + ) + if parsed.empty: + continue + + if config["aggregate_daily"]: + parsed[constants.TIME_INDEX] = parsed[constants.TIME_INDEX].dt.floor("D") + parsed = parsed.groupby(constants.TIME_INDEX, as_index=False)["value"].mean() + + parsed["group_rank"] = payload.get("group_rank", 999) + parsed["series_rank"] = payload.get("series_rank", 999) + frames.append(parsed) + + if not frames: + return self._empty_data_frame(variable) + + combined = pd.concat(frames, ignore_index=True) + combined = combined.sort_values([constants.TIME_INDEX, "group_rank", "series_rank"]) + combined = combined.drop_duplicates(subset=[constants.TIME_INDEX], keep="first") + combined = combined[[constants.TIME_INDEX, "value"]].rename(columns={"value": variable}) + combined = combined.dropna(subset=[variable]).sort_values(constants.TIME_INDEX) + return combined.set_index(constants.TIME_INDEX) + + def get_data( + self, + gauge_id: str, + variable: str, + start_date: Optional[str] = None, + end_date: Optional[str] = None, + ) -> pd.DataFrame: + start_date = utils.format_start_date(start_date) + end_date = utils.format_end_date(end_date) + gauge_id = str(gauge_id).strip() + + if variable not in self.get_available_variables(): + raise ValueError(f"Unsupported variable: {variable}") + + try: + raw_data = self._download_data(gauge_id, variable, start_date, end_date) + df = self._parse_data(gauge_id, raw_data, variable) + except Exception as exc: + logger.error(f"Failed to get data for site {gauge_id}, variable {variable}: {exc}") + return self._empty_data_frame(variable) + + if df.empty: + return df + + start_dt = pd.to_datetime(start_date) + end_dt = pd.to_datetime(end_date) + if "instantaneous" in variable or "hourly" in variable: + end_dt = end_dt + pd.Timedelta(days=1) + return df[(df.index >= start_dt) & (df.index < end_dt)] + + return df[(df.index >= start_dt) & (df.index <= end_dt)] diff --git a/tests/test_data/greece_discharge_checked_data_sample.csv b/tests/test_data/greece_discharge_checked_data_sample.csv new file mode 100644 index 0000000..fb809f9 --- /dev/null +++ b/tests/test_data/greece_discharge_checked_data_sample.csv @@ -0,0 +1,2 @@ +2025-01-01 00:00,0.51, +2025-01-01 01:00,0.52, diff --git a/tests/test_data/greece_discharge_search_sample.json b/tests/test_data/greece_discharge_search_sample.json new file mode 100644 index 0000000..f83696b --- /dev/null +++ b/tests/test_data/greece_discharge_search_sample.json @@ -0,0 +1,43 @@ +{ + "count": 2, + "next": null, + "previous": null, + "results": [ + { + "id": 1458, + "last_update": "2026-03-17T22:00:00Z", + "last_modified": "2018-06-05T22:48:12.746929+03:00", + "name": "Ανθήλη", + "code": "", + "remarks": "", + "geom": "SRID=4326;POINT (22.466853 38.856109)", + "display_timezone": "Etc/GMT-2", + "altitude": null, + "start_date": null, + "end_date": null, + "overseer": "", + "owner": 11 + }, + { + "id": 28082, + "last_update": null, + "last_modified": "2025-07-02T15:11:57.302632+03:00", + "name": "Γούρια (παροχή)", + "code": "", + "remarks": "", + "geom": "SRID=4326;POINT (21.26099 38.48039)", + "display_timezone": "Etc/GMT-2", + "altitude": null, + "start_date": null, + "end_date": null, + "overseer": "", + "owner": 11 + } + ], + "bounding_box": [ + 21.26099, + 38.48039, + 22.466853, + 38.856109 + ] +} diff --git a/tests/test_data/greece_discharge_timeseries_sample.json b/tests/test_data/greece_discharge_timeseries_sample.json new file mode 100644 index 0000000..3ca6d8d --- /dev/null +++ b/tests/test_data/greece_discharge_timeseries_sample.json @@ -0,0 +1,25 @@ +{ + "count": 2, + "next": null, + "previous": null, + "results": [ + { + "id": 9624, + "type": "Initial", + "last_modified": "2018-06-05T23:14:51.111152+03:00", + "time_step": "15min", + "name": "", + "publicly_available": true, + "timeseries_group": 254 + }, + { + "id": 9820, + "type": "Checked", + "last_modified": "2019-12-24T22:43:59.002790+02:00", + "time_step": "15min", + "name": "", + "publicly_available": true, + "timeseries_group": 254 + } + ] +} diff --git a/tests/test_data/greece_level_search_sample.json b/tests/test_data/greece_level_search_sample.json new file mode 100644 index 0000000..f7e1d27 --- /dev/null +++ b/tests/test_data/greece_level_search_sample.json @@ -0,0 +1,28 @@ +{ + "count": 1, + "next": null, + "previous": null, + "results": [ + { + "id": 8424, + "last_update": "2026-03-04T07:45:00Z", + "last_modified": "2021-11-16T13:35:44.555169+02:00", + "name": "Μάνδρα Ρέμα", + "code": "", + "remarks": "", + "geom": "SRID=4326;POINT (23.456293 38.098118)", + "display_timezone": "Etc/GMT-2", + "altitude": 215.0, + "start_date": null, + "end_date": null, + "overseer": "", + "owner": 25 + } + ], + "bounding_box": [ + 23.456293, + 38.098118, + 23.456293, + 38.098118 + ] +} diff --git a/tests/test_data/greece_stage_data_sample.csv b/tests/test_data/greece_stage_data_sample.csv new file mode 100644 index 0000000..3dee5e1 --- /dev/null +++ b/tests/test_data/greece_stage_data_sample.csv @@ -0,0 +1,4 @@ +2025-01-01 00:00,120, +2025-01-01 12:00,130, +2025-01-01 18:00,-999, +2025-01-02 00:00,140, diff --git a/tests/test_data/greece_stage_search_page1_sample.json b/tests/test_data/greece_stage_search_page1_sample.json new file mode 100644 index 0000000..25a2b9b --- /dev/null +++ b/tests/test_data/greece_stage_search_page1_sample.json @@ -0,0 +1,43 @@ +{ + "count": 3, + "next": "https://system.openhi.net/api/stations/?page=2&q=ts_only%3A+variable%3Astage", + "previous": null, + "results": [ + { + "id": 1458, + "last_update": "2026-03-17T22:00:00Z", + "last_modified": "2018-06-05T22:48:12.746929+03:00", + "name": "Ανθήλη", + "code": "", + "remarks": "", + "geom": "SRID=4326;POINT (22.466853 38.856109)", + "display_timezone": "Etc/GMT-2", + "altitude": null, + "start_date": null, + "end_date": null, + "overseer": "", + "owner": 11 + }, + { + "id": 1356, + "last_update": "2026-03-19T12:00:00Z", + "last_modified": "2013-01-22T16:21:46.034119+02:00", + "name": "Καρβελιώτης - Karveliotis", + "code": "", + "remarks": "HYDRONET station, National Observatory of Athens", + "geom": "SRID=4326;POINT (22.223609 37.073479)", + "display_timezone": "Etc/GMT-2", + "altitude": 598.0, + "start_date": "2011-12-16", + "end_date": null, + "overseer": "", + "owner": 17 + } + ], + "bounding_box": [ + 22.223609, + 37.073479, + 22.466853, + 38.856109 + ] +} diff --git a/tests/test_data/greece_stage_search_page2_sample.json b/tests/test_data/greece_stage_search_page2_sample.json new file mode 100644 index 0000000..702887a --- /dev/null +++ b/tests/test_data/greece_stage_search_page2_sample.json @@ -0,0 +1,28 @@ +{ + "count": 3, + "next": null, + "previous": "https://system.openhi.net/api/stations/?q=ts_only%3A+variable%3Astage", + "results": [ + { + "id": 1534, + "last_update": "2026-03-14T22:00:00Z", + "last_modified": "2019-10-29T11:15:59.063566+02:00", + "name": "Αλαμάνα", + "code": "", + "remarks": "", + "geom": "SRID=4326;POINT (22.4952 38.8125)", + "display_timezone": "Etc/GMT-2", + "altitude": null, + "start_date": null, + "end_date": null, + "overseer": "", + "owner": 11 + } + ], + "bounding_box": [ + 22.223609, + 37.073479, + 22.4952, + 38.856109 + ] +} diff --git a/tests/test_data/greece_station_1458_groups_sample.json b/tests/test_data/greece_station_1458_groups_sample.json new file mode 100644 index 0000000..e918fed --- /dev/null +++ b/tests/test_data/greece_station_1458_groups_sample.json @@ -0,0 +1,18 @@ +{ + "count": 1, + "next": null, + "previous": null, + "results": [ + { + "id": 254, + "last_modified": "2018-06-05T23:14:51.111152+03:00", + "name": "", + "hidden": false, + "precision": 2, + "remarks": "", + "gentity": 1458, + "variable": 2, + "unit_of_measurement": 18 + } + ] +} diff --git a/tests/test_data/greece_station_8424_groups_sample.json b/tests/test_data/greece_station_8424_groups_sample.json new file mode 100644 index 0000000..ecb4083 --- /dev/null +++ b/tests/test_data/greece_station_8424_groups_sample.json @@ -0,0 +1,18 @@ +{ + "count": 1, + "next": null, + "previous": null, + "results": [ + { + "id": 881, + "last_modified": "2021-11-19T19:41:37.900822+02:00", + "name": "Water Level", + "hidden": false, + "precision": 2, + "remarks": "", + "gentity": 8424, + "variable": 5710, + "unit_of_measurement": 2 + } + ] +} diff --git a/tests/test_data/greece_water_level_timeseries_sample.json b/tests/test_data/greece_water_level_timeseries_sample.json new file mode 100644 index 0000000..c8ecf0e --- /dev/null +++ b/tests/test_data/greece_water_level_timeseries_sample.json @@ -0,0 +1,16 @@ +{ + "count": 1, + "next": null, + "previous": null, + "results": [ + { + "id": 10306, + "type": "Initial", + "last_modified": "2022-03-23T12:24:26.808693+02:00", + "time_step": "15min", + "name": "", + "publicly_available": true, + "timeseries_group": 881 + } + ] +} diff --git a/tests/test_greece.py b/tests/test_greece.py new file mode 100644 index 0000000..28f9bf1 --- /dev/null +++ b/tests/test_greece.py @@ -0,0 +1,462 @@ +import json +import unittest +from pathlib import Path +from unittest.mock import MagicMock, patch + +import pandas as pd +from pandas.testing import assert_frame_equal + +from rivretrieve import GreeceFetcher, constants + + +class TestGreeceFetcher(unittest.TestCase): + def setUp(self): + self.fetcher = GreeceFetcher() + self.test_data_dir = Path(__file__).parent / "test_data" + + def _load_json(self, filename): + with open(self.test_data_dir / filename, "r", encoding="utf-8") as f: + return json.load(f) + + def _load_text(self, filename): + with open(self.test_data_dir / filename, "r", encoding="utf-8") as f: + return f.read() + + @staticmethod + def _json_response(payload): + response = MagicMock() + response.json.return_value = payload + response.raise_for_status = MagicMock() + return response + + @staticmethod + def _text_response(payload): + response = MagicMock() + response.text = payload + response.raise_for_status = MagicMock() + return response + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_metadata_merges_stage_water_level_and_discharge_queries(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + + stage_page_1 = self._load_json("greece_stage_search_page1_sample.json") + stage_page_2 = self._load_json("greece_stage_search_page2_sample.json") + level_payload = self._load_json("greece_level_search_sample.json") + discharge_payload = self._load_json("greece_discharge_search_sample.json") + + def side_effect(url, params=None, timeout=60): + if url == f"{self.fetcher.BASE_URL}/stations/" and params == {"q": "ts_only: variable:stage"}: + return self._json_response(stage_page_1) + if url == stage_page_1["next"] and params is None: + return self._json_response(stage_page_2) + if url == f"{self.fetcher.BASE_URL}/stations/" and params == {"q": "ts_only: variable:level"}: + return self._json_response(level_payload) + if url == f"{self.fetcher.BASE_URL}/stations/" and params == {"q": "ts_only: variable:discharge"}: + return self._json_response(discharge_payload) + if url == f"{self.fetcher.BASE_URL}/organizations/11/" and params is None: + return self._json_response({"id": 11, "name": "Hellenic Centre for Marine Research"}) + if url == f"{self.fetcher.BASE_URL}/organizations/17/" and params is None: + return self._json_response({"id": 17, "name": "National Observatory of Athens"}) + if url == f"{self.fetcher.BASE_URL}/organizations/25/" and params is None: + return self._json_response({"id": 25, "name": "Mandra Project"}) + raise AssertionError(f"Unexpected request: url={url}, params={params}") + + mock_session.get.side_effect = side_effect + + result_df = self.fetcher.get_metadata() + + self.assertEqual(result_df.index.name, constants.GAUGE_ID) + self.assertEqual(list(result_df.index), ["1356", "1458", "1534", "28082", "8424"]) + self.assertEqual(result_df.loc["1458", constants.STATION_NAME], "Ανθήλη") + self.assertAlmostEqual(result_df.loc["1458", constants.LATITUDE], 38.856109) + self.assertAlmostEqual(result_df.loc["1458", constants.LONGITUDE], 22.466853) + self.assertEqual(result_df.loc["1356", "owner_name"], "National Observatory of Athens") + self.assertEqual(result_df.loc["8424", "owner_name"], "Mandra Project") + self.assertEqual(result_df.loc["28082", constants.SOURCE], self.fetcher.SOURCE) + self.assertEqual(result_df.loc["28082", constants.COUNTRY], "Greece") + self.assertTrue(all(call.kwargs["timeout"] == 60 for call in mock_session.get.call_args_list)) + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_daily_stage_normalizes_water_level_and_converts_cm(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + + groups_payload = self._load_json("greece_station_8424_groups_sample.json") + timeseries_payload = self._load_json("greece_water_level_timeseries_sample.json") + data_payload = self._load_text("greece_stage_data_sample.csv") + + def side_effect(url, params=None, timeout=60): + if url == f"{self.fetcher.BASE_URL}/stations/8424/timeseriesgroups/" and params is None: + return self._json_response(groups_payload) + if url == f"{self.fetcher.BASE_URL}/stations/8424/timeseriesgroups/881/timeseries/" and params is None: + return self._json_response(timeseries_payload) + if url == f"{self.fetcher.BASE_URL}/stations/8424/timeseriesgroups/881/timeseries/10306/data/": + self.assertEqual( + params, + {"fmt": "csv", "start_date": "2025-01-01", "end_date": "2025-01-02"}, + ) + return self._text_response(data_payload) + if url == f"{self.fetcher.BASE_URL}/units/2/" and params is None: + return self._json_response({"id": 2, "symbol": "cm"}) + raise AssertionError(f"Unexpected request: url={url}, params={params}") + + mock_session.get.side_effect = side_effect + + result_df = self.fetcher.get_data( + gauge_id="8424", + variable=constants.STAGE_DAILY_MEAN, + start_date="2025-01-01", + end_date="2025-01-02", + ) + + expected_df = pd.DataFrame( + { + constants.TIME_INDEX: pd.to_datetime(["2025-01-01", "2025-01-02"]), + constants.STAGE_DAILY_MEAN: [1.25, 1.40], + } + ).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + self.assertEqual(result_df.index.name, constants.TIME_INDEX) + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_instant_discharge_prefers_initial_subdaily_series(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + + groups_payload = { + "count": 1, + "next": None, + "previous": None, + "results": [ + { + "id": 6, + "last_modified": "2019-12-24T22:43:59.004465+02:00", + "name": "", + "hidden": False, + "precision": 2, + "remarks": "", + "gentity": 1458, + "variable": 2, + "unit_of_measurement": 18, + } + ], + } + timeseries_payload = { + "count": 3, + "next": None, + "previous": None, + "results": [ + { + "id": 10105, + "type": "Initial", + "last_modified": "2020-11-18T16:36:51.279025+02:00", + "time_step": "", + "name": "", + "publicly_available": True, + "timeseries_group": 6, + }, + { + "id": 10106, + "type": "Aggregated", + "last_modified": "2020-11-18T16:38:15.700871+02:00", + "time_step": "1h", + "name": "Mean", + "publicly_available": True, + "timeseries_group": 6, + }, + { + "id": 10115, + "type": "Aggregated", + "last_modified": "2020-12-01T21:57:52.001775+02:00", + "time_step": "1D", + "name": "Mean", + "publicly_available": True, + "timeseries_group": 6, + }, + ], + } + data_payload = "2025-01-01 00:00,30.44,\n2025-01-01 01:00,30.77,\n" + + def side_effect(url, params=None, timeout=60): + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/" and params is None: + return self._json_response(groups_payload) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/" and params is None: + return self._json_response(timeseries_payload) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10105/data/": + return self._text_response(data_payload) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10106/data/": + return self._text_response("") + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10115/data/": + return self._text_response("") + if url == f"{self.fetcher.BASE_URL}/units/18/" and params is None: + return self._json_response({"id": 18, "symbol": "m³/s"}) + raise AssertionError(f"Unexpected request: url={url}, params={params}") + + mock_session.get.side_effect = side_effect + + result_df = self.fetcher.get_data( + gauge_id="1458", + variable=constants.DISCHARGE_INSTANT, + start_date="2025-01-01", + end_date="2025-01-01", + ) + + expected_df = pd.DataFrame( + { + constants.TIME_INDEX: pd.to_datetime(["2025-01-01 00:00:00", "2025-01-01 01:00:00"]), + constants.DISCHARGE_INSTANT: [30.44, 30.77], + } + ).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + self.assertEqual(result_df.index.name, constants.TIME_INDEX) + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_daily_stage_falls_back_when_primary_group_is_empty(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + + groups_payload = { + "count": 2, + "next": None, + "previous": None, + "results": [ + { + "id": 12, + "last_modified": "2025-10-22T16:19:48.103896+03:00", + "name": "Mean stage", + "hidden": False, + "precision": 2, + "remarks": "", + "gentity": 9999, + "variable": 14, + "unit_of_measurement": 6, + }, + { + "id": 13, + "last_modified": "2024-10-22T16:19:48.103896+03:00", + "name": "Stage", + "hidden": False, + "precision": 2, + "remarks": "", + "gentity": 9999, + "variable": 14, + "unit_of_measurement": 6, + }, + ], + } + primary_timeseries_payload = { + "count": 1, + "next": None, + "previous": None, + "results": [ + { + "id": 200, + "type": "Aggregated", + "last_modified": "2025-10-22T16:19:48.103896+03:00", + "time_step": "1D", + "name": "Mean", + "publicly_available": True, + "timeseries_group": 12, + } + ], + } + fallback_timeseries_payload = { + "count": 1, + "next": None, + "previous": None, + "results": [ + { + "id": 201, + "type": "Initial", + "last_modified": "2024-10-22T16:19:48.103896+03:00", + "time_step": "15min", + "name": "", + "publicly_available": True, + "timeseries_group": 13, + } + ], + } + + def side_effect(url, params=None, timeout=60): + if url == f"{self.fetcher.BASE_URL}/stations/9999/timeseriesgroups/" and params is None: + return self._json_response(groups_payload) + if url == f"{self.fetcher.BASE_URL}/stations/9999/timeseriesgroups/12/timeseries/" and params is None: + return self._json_response(primary_timeseries_payload) + if url == f"{self.fetcher.BASE_URL}/stations/9999/timeseriesgroups/13/timeseries/" and params is None: + return self._json_response(fallback_timeseries_payload) + if url == f"{self.fetcher.BASE_URL}/stations/9999/timeseriesgroups/12/timeseries/200/data/": + return self._text_response("") + if url == f"{self.fetcher.BASE_URL}/stations/9999/timeseriesgroups/13/timeseries/201/data/": + return self._text_response("2025-01-01 00:00,1.00,\n2025-01-01 12:00,3.00,\n") + if url == f"{self.fetcher.BASE_URL}/units/6/" and params is None: + return self._json_response({"id": 6, "symbol": "m"}) + raise AssertionError(f"Unexpected request: url={url}, params={params}") + + mock_session.get.side_effect = side_effect + + result_df = self.fetcher.get_data( + gauge_id="9999", + variable=constants.STAGE_DAILY_MEAN, + start_date="2025-01-01", + end_date="2025-01-01", + ) + + expected_df = pd.DataFrame( + { + constants.TIME_INDEX: pd.to_datetime(["2025-01-01"]), + constants.STAGE_DAILY_MEAN: [2.0], + } + ).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_daily_discharge_falls_back_from_empty_aggregated_series_to_initial(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + + groups_payload = { + "count": 1, + "next": None, + "previous": None, + "results": [ + { + "id": 6, + "last_modified": "2020-12-01T21:57:52.001775+02:00", + "name": "", + "hidden": False, + "precision": 2, + "remarks": "", + "gentity": 1458, + "variable": 2, + "unit_of_measurement": 18, + } + ], + } + timeseries_payload = { + "count": 3, + "next": None, + "previous": None, + "results": [ + { + "id": 10105, + "type": "Initial", + "last_modified": "2020-11-18T16:36:51.279025+02:00", + "time_step": "", + "name": "", + "publicly_available": True, + "timeseries_group": 6, + }, + { + "id": 10106, + "type": "Aggregated", + "last_modified": "2020-11-18T16:38:15.700871+02:00", + "time_step": "1h", + "name": "Mean", + "publicly_available": True, + "timeseries_group": 6, + }, + { + "id": 10115, + "type": "Aggregated", + "last_modified": "2020-12-01T21:57:52.001775+02:00", + "time_step": "1D", + "name": "Mean", + "publicly_available": True, + "timeseries_group": 6, + }, + ], + } + + def side_effect(url, params=None, timeout=60): + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/" and params is None: + return self._json_response(groups_payload) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/" and params is None: + return self._json_response(timeseries_payload) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10115/data/": + return self._text_response("") + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10105/data/": + return self._text_response( + "2024-04-14 07:01,15.30,\n" + "2024-04-14 08:01,15.31,\n" + "2024-04-15 07:01,14.30,\n" + "2024-04-15 08:01,14.70,\n" + ) + if url == f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/6/timeseries/10106/data/": + return self._text_response("") + if url == f"{self.fetcher.BASE_URL}/units/18/" and params is None: + return self._json_response({"id": 18, "symbol": "m³/s"}) + raise AssertionError(f"Unexpected request: url={url}, params={params}") + + mock_session.get.side_effect = side_effect + + result_df = self.fetcher.get_data( + gauge_id="1458", + variable=constants.DISCHARGE_DAILY_MEAN, + start_date="2024-04-14", + end_date="2024-04-15", + ) + + expected_df = pd.DataFrame( + { + constants.TIME_INDEX: pd.to_datetime(["2024-04-14", "2024-04-15"]), + constants.DISCHARGE_DAILY_MEAN: [15.305, 14.5], + } + ).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + self.assertEqual(result_df.index.name, constants.TIME_INDEX) + + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_returns_standardized_empty_frame_when_no_groups_match(self, mock_requests_session): + mock_session = MagicMock() + mock_requests_session.return_value = mock_session + self.fetcher = GreeceFetcher() + mock_session.get.return_value = self._json_response({"count": 0, "next": None, "previous": None, "results": []}) + + result_df = self.fetcher.get_data( + gauge_id="1458", + variable=constants.STAGE_INSTANT, + start_date="2025-01-01", + end_date="2025-01-02", + ) + + expected_df = pd.DataFrame(columns=[constants.TIME_INDEX, constants.STAGE_INSTANT]).set_index( + constants.TIME_INDEX + ) + + assert_frame_equal(result_df, expected_df) + self.assertEqual(result_df.index.name, constants.TIME_INDEX) + mock_session.get.assert_called_once_with( + f"{self.fetcher.BASE_URL}/stations/1458/timeseriesgroups/", + params=None, + timeout=60, + ) + + def test_available_variables(self): + self.assertEqual( + self.fetcher.get_available_variables(), + ( + constants.DISCHARGE_DAILY_MEAN, + constants.DISCHARGE_INSTANT, + constants.STAGE_DAILY_MEAN, + constants.STAGE_INSTANT, + ), + ) + + def test_unsupported_variable_raises(self): + with self.assertRaises(ValueError): + self.fetcher.get_data("1458", constants.WATER_TEMPERATURE_DAILY_MEAN) + + +if __name__ == "__main__": + unittest.main()