diff --git a/.gitignore b/.gitignore index fd149f7..a12902b 100644 --- a/.gitignore +++ b/.gitignore @@ -8,4 +8,5 @@ downloaded_data/ rivretrieve/data **/*.pyc -**/*.png \ No newline at end of file +**/*.png +**/*.env \ No newline at end of file diff --git a/examples/test_brazil_fetcher.py b/examples/test_brazil_fetcher.py index 186fe2b..08e55fa 100644 --- a/examples/test_brazil_fetcher.py +++ b/examples/test_brazil_fetcher.py @@ -7,23 +7,27 @@ logging.basicConfig(level=logging.INFO) gauge_ids = [ - "12650000", + "15400000", ] -variable = constants.DISCHARGE +variable = constants.STAGE_DAILY_MEAN +start_date = "2023-01-01" +end_date = "2024-12-31" plt.figure(figsize=(12, 6)) +# Initialize fetcher (credentials are loaded from .env by default) fetcher = BrazilFetcher() + for gauge_id in gauge_ids: - print(f"Fetching data for {gauge_id}...") - data = fetcher.get_data(gauge_id=gauge_id, variable=variable) + print(f"Fetching {variable} for {gauge_id} from {start_date} to {end_date}...") + data = fetcher.get_data(gauge_id=gauge_id, variable=variable, start_date=start_date, end_date=end_date) if not data.empty: print(f"Data for {gauge_id}:") print(data.head()) - print(f"Time series from {data[constants.TIME_INDEX].min()} to {data[constants.TIME_INDEX].max()}") + print(f"Time series from {data.index.min()} to {data.index.max()}") plt.plot( - data[constants.TIME_INDEX], - data[constants.DISCHARGE], + data.index, + data[variable], label=gauge_id, marker=".", linestyle="-", @@ -33,8 +37,8 @@ if "data" in locals() and not data.empty: plt.xlabel(constants.TIME_INDEX) - plt.ylabel(f"{constants.DISCHARGE} (m3/s)") - plt.title(f"Brazil River Discharge ({gauge_ids[0]} - Full Time Series)") + plt.ylabel(f"{variable} (m3/s)") + plt.title(f"Brazil River Discharge ({gauge_ids[0]} - {start_date} to {end_date})") plt.legend() plt.grid(True) plt.tight_layout() diff --git a/requirements.txt b/requirements.txt index 7ca646e..75e242b 100644 --- a/requirements.txt +++ b/requirements.txt @@ -8,5 +8,7 @@ dataretrieval>=1.0.0 openpyxl>=3.0.0 xarray parameterized +python-dotenv +ruff tqdm zarr>=3.0.7 \ No newline at end of file diff --git a/rivretrieve/brazil.py b/rivretrieve/brazil.py index a9e3df4..b92da1f 100644 --- a/rivretrieve/brazil.py +++ b/rivretrieve/brazil.py @@ -1,23 +1,45 @@ """Fetcher for Brazilian river gauge data from ANA Hidroweb.""" -import io import logging -import re -import zipfile -from typing import Optional +import os +import time +from datetime import datetime +from typing import Any, Dict, List, Optional import pandas as pd import requests +from dotenv import load_dotenv from . import base, constants, utils logger = logging.getLogger(__name__) +# Load environment variables from .env file +load_dotenv(dotenv_path=os.path.join((os.path.dirname(__file__)), ".env")) + +# Get credentials from environment variables +USERNAME = os.environ.get("ANA_USERNAME") +PASSWORD = os.environ.get("ANA_PASSWORD") -class BrazilFetcher(base.RiverDataFetcher): - """Fetches river gauge data from Brazil's ANA Hidroweb.""" - BASE_URL = "https://www.snirh.gov.br/hidroweb/rest/api/documento/convencionais" +class BrazilFetcher(base.RiverDataFetcher): + """Fetches river gauge data from Brazil's ANA Hidroweb API v2.""" + + BASE_URL = "https://www.ana.gov.br/hidrowebservice/EstacoesTelemetricas" + AUTH_URL = f"{BASE_URL}/OAUth/v1" + + def __init__(self, username: Optional[str] = None, password: Optional[str] = None): + super().__init__() + self.username = username or USERNAME + self.password = password or PASSWORD + self._token = None + self._token_expiry = 0 + + if not self.username or not self.password: + logger.error( + "ANA Username or Password not provided. Please set ANA_USERNAME and ANA_PASSWORD in ," + "your .env file or pass them to the constructor." + ) @staticmethod def get_gauge_ids() -> pd.DataFrame: @@ -26,124 +48,249 @@ def get_gauge_ids() -> pd.DataFrame: @staticmethod def get_available_variables() -> tuple[str, ...]: - return (constants.DISCHARGE, constants.STAGE) - - def _download_data(self, gauge_id: str, variable: str, start_date: str, end_date: str) -> Optional[pd.DataFrame]: - """Downloads and extracts the data file.""" - params = {"tipo": 3, "documentos": gauge_id} + return (constants.DISCHARGE_DAILY_MEAN, constants.STAGE_DAILY_MEAN) + + def get_metadata(self) -> pd.DataFrame: + """Fetches station metadata for all Brazilian states.""" + if not self.username or not self.password: + logger.error("ANA Username or Password not provided.") + return pd.DataFrame().set_index(constants.GAUGE_ID) + + token = self._get_token() + if not token: + logger.error("Cannot fetch metadata without a token.") + return pd.DataFrame().set_index(constants.GAUGE_ID) + + states = [ + "AC", + "AL", + "AM", + "AP", + "BA", + "CE", + "DF", + "ES", + "GO", + "MA", + "MT", + "MS", + "MG", + "PA", + "PB", + "PR", + "PE", + "PI", + "RJ", + "RN", + "RS", + "RO", + "RR", + "SC", + "SP", + "SE", + "TO", + ] + + metadata_url = f"{self.BASE_URL}/HidroInventarioEstacoes/v1" + all_stations = [] s = utils.requests_retry_session() + headers = {"Authorization": f"Bearer {token}"} + + for state in states: + logger.info(f"Fetching metadata for state: {state}") + params = {"Unidade Federativa": state} + try: + response = s.get(metadata_url, params=params, headers=headers) + response.raise_for_status() + data = response.json() + if isinstance(data, list): + all_stations.extend(data) + elif isinstance(data, dict) and data.get("status") == "OK" and data.get("items"): + all_stations.extend(data["items"]) + else: + logger.warning(f"No stations found for state {state} or unexpected response: {data}") + except requests.exceptions.RequestException as e: + logger.error(f"Error fetching metadata for state {state}: {e}") + except Exception as e: + logger.error(f"Error processing metadata for state {state}: {e}") + time.sleep(0.1) + + if not all_stations: + return pd.DataFrame().set_index(constants.GAUGE_ID) + + df = pd.DataFrame(all_stations) + + rename_map = { + "codigoestacao": constants.GAUGE_ID, + "Estacao_Nome": constants.STATION_NAME, + "Latitude": constants.LATITUDE, + "Longitude": constants.LONGITUDE, + "Altitude": constants.ALTITUDE, + "Area_Drenagem": constants.AREA, + "Bacia_Nome": constants.RIVER, + } + df = df.rename(columns=rename_map) + df[constants.GAUGE_ID] = df[constants.GAUGE_ID].astype(str) + return df.set_index(constants.GAUGE_ID) + + def _get_token(self) -> Optional[str]: + """Gets and caches the authentication token.""" + if not self.username or not self.password: + return None + if self._token and time.time() < self._token_expiry: + return self._token + + logger.info("Fetching new authentication token for Brazil...") + headers = {"accept": "*/*", "Identificador": self.username, "Senha": self.password} + s = utils.requests_retry_session() try: - response = s.get(self.BASE_URL, params=params, headers=utils.DEFAULT_HEADERS) + response = s.get(self.AUTH_URL, headers=headers) response.raise_for_status() - - with zipfile.ZipFile(io.BytesIO(response.content)) as zf: - # Find the correct inner zip file (vazoes for discharge, cotas for stage) - inner_zip_name = None - pattern = f"^{gauge_id}/(vazoes|cotas)_{gauge_id}.zip$" - for name in zf.namelist(): - if re.match(pattern, name): - if (variable == constants.DISCHARGE and "vazoes" in name) or ( - variable == constants.STAGE and "cotas" in name - ): - inner_zip_name = name - break - - if not inner_zip_name: - logger.warning(f"Could not find inner zip for site {gauge_id}, variable {variable}") - return None - - with zf.open(inner_zip_name) as inner_zf_file: - with zipfile.ZipFile(io.BytesIO(inner_zf_file.read())) as inner_zf: - data_file_name = inner_zf.namelist()[0] - with inner_zf.open(data_file_name) as data_file: - # The files are ISO-8859-1 encoded - raw_df = pd.read_csv( - data_file, - sep=";", - decimal=",", - skiprows=13, - encoding="ISO-8859-1", - engine="python", - ) - return raw_df - + data = response.json() + if data.get("status") == "OK" and data.get("items", {}).get("sucesso"): + self._token = data["items"]["tokenautenticacao"] + # Set expiry to 14 minutes (840 seconds) to be safe + self._token_expiry = time.time() + 840 + logger.info("Successfully obtained new token.") + return self._token + else: + logger.error(f"Authentication failed: {data}") + return None except requests.exceptions.RequestException as e: - logger.error(f"Error downloading from Hidroweb for site {gauge_id}: {e}") - return None - except zipfile.BadZipFile: - logger.error(f"Bad zip file for site {gauge_id}") + logger.error(f"Error fetching token: {e}") return None except Exception as e: - logger.error(f"Error processing file for site {gauge_id}: {e}") + logger.error(f"Error processing token response: {e}") return None - def _parse_data(self, gauge_id: str, raw_df: Optional[pd.DataFrame], variable: str) -> pd.DataFrame: - """Parses the raw DataFrame.""" - if raw_df is None or raw_df.empty: - return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) + def _download_data(self, gauge_id: str, variable: str, start_date: str, end_date: str) -> List[Dict[str, Any]]: + """Downloads raw data in yearly chunks.""" + if not self.username or not self.password: + return [] + + if variable == constants.DISCHARGE_DAILY_MEAN: + data_url = f"{self.BASE_URL}/HidroSerieVazao/v1" + elif variable == constants.STAGE_DAILY_MEAN: + data_url = f"{self.BASE_URL}/HidroSerieCotas/v1" + else: + logger.error(f"Unsupported variable for daily download: {variable}") + return [] + + all_data = [] + start_dt = datetime.strptime(start_date, "%Y-%m-%d") + end_dt = datetime.strptime(end_date, "%Y-%m-%d") + s = utils.requests_retry_session() - try: - id_cols = ["EstacaoCodigo", "NivelConsistencia", "Data", "Hora"] - if variable == constants.DISCHARGE: - prefix = "Vazao" - id_cols.append("MetodoObtencaoVazoes") - else: # stage - prefix = "Cota" - - # Select relevant columns - value_cols = [col for col in raw_df.columns if col.startswith(prefix)] - status_cols = [col for col in value_cols if col.endswith("Status")] - value_cols = [col for col in value_cols if not col.endswith("Status")] - - df = raw_df[id_cols + value_cols + status_cols].copy() - - # Rename status columns for pivot - rename_map = {sc: sc.replace("Status", "_Status") for sc in status_cols} - df = df.rename(columns=rename_map) - - # Rename value columns for pivot - rename_map = {vc: f"{vc}_Value" for vc in value_cols} - df = df.rename(columns=rename_map) - - # Pivot longer - df_long = pd.melt(df, id_vars=id_cols, var_name="day_type", value_name="Value") - df_long[["day", "Type"]] = df_long["day_type"].str.split("_", expand=True) - df_long["day"] = df_long["day"].str.replace(prefix, "").astype(int) - - # Separate Value and Status - df_values = df_long[df_long["Type"] == "Value"].copy() - df_status = df_long[df_long["Type"] == "Status"].copy() - df_status = df_status.rename(columns={"Value": "Status"}) - - # Merge back Value and Status - id_cols_melt = id_cols + ["day"] - df_merged = pd.merge( - df_values[id_cols_melt + ["Value"]], - df_status[id_cols_melt + ["Status"]], - on=id_cols_melt, - how="left", + current_year = start_dt.year + while current_year <= end_dt.year: + token = self._get_token() + if not token: + logger.error("Cannot download data without a token.") + break + + year_start = datetime(current_year, 1, 1) + year_end = datetime(current_year, 12, 31) + + req_start_date = max(start_dt, year_start) + req_end_date = min(end_dt, year_end) + + if req_start_date > req_end_date: + current_year += 1 + continue + + # Manually build the URL to control encoding + base_data_url = f"{data_url}?" + params_list = [ + f"C%C3%B3digo%20da%20Esta%C3%A7%C3%A3o={gauge_id}", + "Tipo%20Filtro%20Data=DATA_LEITURA", + f"Data%20Inicial%20(yyyy-MM-dd)={req_start_date.strftime('%Y-%m-%d')}", + f"Data%20Final%20(yyyy-MM-dd)={req_end_date.strftime('%Y-%m-%d')}", + ] + full_url = base_data_url + "&".join(params_list) + + headers = {"accept": "*/*", "Authorization": f"Bearer {token}"} + + logger.debug( + f"Fetching {variable} for site {gauge_id} from {req_start_date.strftime('%Y-%m-%d')} to " + f" {req_end_date.strftime('%Y-%m-%d')}" ) + logger.debug(f"Request URL: {full_url}") + try: + response = s.get(full_url, headers=headers) + response.raise_for_status() + data = response.json() + if isinstance(data, list): + all_data.extend(data) + elif isinstance(data, dict) and data.get("status") == "OK" and data.get("items"): + all_data.extend(data["items"]) + elif isinstance(data, dict) and data.get("status") == "OK": + logger.info(f"No items returned for {gauge_id} for year {current_year}") + else: + logger.warning(f"API returned unexpected response: {data}") + + except requests.exceptions.RequestException as e: + logger.error(f"Error downloading data chunk for {gauge_id}: {e}") + except Exception as e: + logger.error(f"Error processing data chunk for {gauge_id}: {e}") + + current_year += 1 + time.sleep(0.2) # Be nice to the API + + return all_data + + def _parse_data(self, gauge_id: str, raw_data: List[Dict[str, Any]], variable: str) -> pd.DataFrame: + """Parses the raw JSON data from daily endpoints.""" + if not raw_data: + return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) - # Create Date column - df_merged["Data"] = pd.to_datetime(df_merged["Data"], format="%d/%m/%Y") - df_merged["year"] = df_merged["Data"].dt.year - df_merged["month"] = df_merged["Data"].dt.month - - # Handle potential errors in make_date by coercing - date_df = df_merged[["year", "month", "day"]].copy() - df_merged[constants.TIME_INDEX] = pd.to_datetime(date_df, errors="coerce") - df_merged = df_merged.dropna(subset=[constants.TIME_INDEX]) - - df_merged["Value"] = pd.to_numeric(df_merged["Value"], errors="coerce") - if variable == constants.STAGE: # cm to m - df_merged["Value"] = df_merged["Value"] / 100.0 - - df_final = df_merged[[constants.TIME_INDEX, "Value"]].rename(columns={"Value": variable}) - return df_final.dropna().sort_values(by=constants.TIME_INDEX).reset_index(drop=True) - + all_dfs = [] + try: + for month_data in raw_data: + if not isinstance(month_data, dict): + logger.warning(f"Unexpected item format in raw_data: {month_data}") + continue + + month_str = month_data.get("Data_Hora_Dado") + if not month_str: + logger.warning(f"Missing 'Data_Hora_Dado' in month_data: {month_data}") + continue + + year = int(month_str[:4]) + month = int(month_str[5:7]) + + if variable == constants.DISCHARGE_DAILY_MEAN: + val_prefix = "Vazao_" + unit_conversion = 1.0 + elif variable == constants.STAGE_DAILY_MEAN: + val_prefix = "Cota_" + unit_conversion = 0.01 # cm to m + else: + return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) + + day_values = {} + for day in range(1, 32): + day_str = f"{day:02d}" + val_col = f"{val_prefix}{day_str}" + if val_col in month_data and month_data[val_col] is not None: + try: + date = datetime(year, month, day) + day_values[date] = pd.to_numeric(month_data[val_col], errors="coerce") * unit_conversion + except ValueError: + # Handles invalid dates like Feb 30 + continue + + month_df = pd.DataFrame(list(day_values.items()), columns=[constants.TIME_INDEX, variable]) + all_dfs.append(month_df) + + if not all_dfs: + return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) + + df = pd.concat(all_dfs, ignore_index=True) + df = df.dropna().sort_values(by=constants.TIME_INDEX).set_index(constants.TIME_INDEX) + return df except Exception as e: - logger.error(f"Error parsing data for site {gauge_id}: {e}") + logger.error(f"Error parsing CSV data for site {gauge_id}: {e}") return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) def get_data( @@ -154,6 +301,10 @@ def get_data( end_date: Optional[str] = None, ) -> pd.DataFrame: """Fetches and parses Brazilian river gauge data.""" + if not self.username or not self.password: + logger.error("ANA Username or Password not provided. Check your .env file or constructor arguments.") + return pd.DataFrame(columns=[constants.TIME_INDEX, variable]) + start_date = utils.format_start_date(start_date) end_date = utils.format_end_date(end_date) if variable not in self.get_available_variables(): @@ -162,11 +313,9 @@ def get_data( try: raw_data = self._download_data(gauge_id, variable, start_date, end_date) df = self._parse_data(gauge_id, raw_data, variable) - - # Filter by date range start_date_dt = pd.to_datetime(start_date) end_date_dt = pd.to_datetime(end_date) - df = df[(df[constants.TIME_INDEX] >= start_date_dt) & (df[constants.TIME_INDEX] <= end_date_dt)] + df = df[(df.index >= start_date_dt) & (df.index <= end_date_dt)] return df except Exception as e: logger.error(f"Failed to get data for site {gauge_id}, variable {variable}: {e}") diff --git a/tests/test_brazil.py b/tests/test_brazil.py new file mode 100644 index 0000000..dc92dc4 --- /dev/null +++ b/tests/test_brazil.py @@ -0,0 +1,85 @@ +import unittest +from unittest.mock import MagicMock, patch + +import pandas as pd +from pandas.testing import assert_frame_equal + +from rivretrieve import BrazilFetcher, constants + + +class TestBrazilFetcher(unittest.TestCase): + def setUp(self): + self.fetcher = BrazilFetcher(username="testuser", password="testpass") + + @patch("rivretrieve.brazil.BrazilFetcher._get_token") + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_discharge(self, mock_session, mock_get_token): + mock_get_token.return_value = "fake_token" + mock_response = MagicMock() + mock_json = [ + { + "Data_Hora_Dado": "2024-01-01 00:00:00.0", + "Vazao_01": "10.0", + "Vazao_02": "11.0", + "Vazao_03": "12.0", + # ... add other days up to 31, some can be None + "Vazao_31": None, + } + ] + mock_response.json.return_value = mock_json + mock_response.raise_for_status = MagicMock() + mock_session.return_value.get.return_value = mock_response + + gauge_id = "12345678" + variable = constants.DISCHARGE_DAILY_MEAN + start_date = "2024-01-01" + end_date = "2024-01-03" + + result_df = self.fetcher.get_data(gauge_id, variable, start_date, end_date) + + expected_data = { + constants.TIME_INDEX: pd.to_datetime(["2024-01-01", "2024-01-02", "2024-01-03"]), + constants.DISCHARGE_DAILY_MEAN: [10.0, 11.0, 12.0], + } + expected_df = pd.DataFrame(expected_data).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + mock_session.return_value.get.assert_called_once() + + @patch("rivretrieve.brazil.BrazilFetcher._get_token") + @patch("rivretrieve.utils.requests_retry_session") + def test_get_data_stage(self, mock_session, mock_get_token): + mock_get_token.return_value = "fake_token" + mock_response = MagicMock() + mock_json = [ + { + "Data_Hora_Dado": "2024-02-01 00:00:00.0", + "Cota_01": "150", + "Cota_02": "155", + # ... add other days up to 29 + "Cota_29": "160", + } + ] + mock_response.json.return_value = mock_json + mock_response.raise_for_status = MagicMock() + mock_session.return_value.get.return_value = mock_response + + gauge_id = "12345678" + variable = constants.STAGE_DAILY_MEAN + start_date = "2024-02-01" + end_date = "2024-02-02" + + result_df = self.fetcher.get_data(gauge_id, variable, start_date, end_date) + + expected_data = { + constants.TIME_INDEX: pd.to_datetime(["2024-02-01", "2024-02-02"]), + constants.STAGE_DAILY_MEAN: [1.50, 1.55], # Converted to meters + } + expected_df = pd.DataFrame(expected_data).set_index(constants.TIME_INDEX) + + assert_frame_equal(result_df, expected_df) + mock_session.return_value.get.assert_called_once() + + +if __name__ == "__main__": + unittest.main()