Skip to content
2 changes: 2 additions & 0 deletions src/cosecha/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from cosecha import exceptions
from cosecha._logging import configure_logger
from cosecha.reaping import (
ASOSReaper,
GriddedReaper,
MRMSReaper,
NWPReaper,
Expand All @@ -20,6 +21,7 @@
__version__ = "999"

__all__ = [
"ASOSReaper",
"GriddedReaper",
"MRMSReaper",
"NWPReaper",
Expand Down
2 changes: 2 additions & 0 deletions src/cosecha/reaping/__init__.py
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,15 @@

from __future__ import annotations

from cosecha.reaping.asos import ASOSReaper
from cosecha.reaping.base import GriddedReaper, TimeSeriesReaper
from cosecha.reaping.mrms import MRMSReaper
from cosecha.reaping.nwis import USGSNWISReaper
from cosecha.reaping.nwp import NWPReaper
from cosecha.reaping.usace import ReservoirReaper

__all__ = [
"ASOSReaper",
"GriddedReaper",
"MRMSReaper",
"NWPReaper",
Expand Down
150 changes: 150 additions & 0 deletions src/cosecha/reaping/asos.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
"""IEM ASOS data reapers.

This module provides implementations for harvesting ASOS observations
from the Iowa Environmental Mesonet (IEM) API.
"""

from __future__ import annotations

from io import StringIO
from typing import Any
from urllib.parse import urlencode

import pandas as pd
import tiny_retriever

from cosecha._logging import logger
from cosecha._utils import apply_ts_transformations, wrap_errors
from cosecha.exceptions import APIError, DateRangeError, InvalidSiteError
from cosecha.reaping.base import TimeSeriesReaper

__all__ = [
"ASOSReaper",
]

BASE_URL = "https://mesonet.agron.iastate.edu/cgi-bin/request/asos.py"


class ASOSReaper(TimeSeriesReaper):
"""Reaper for IEM ASOS data."""

timeout: int = 120
Comment thread
sray014 marked this conversation as resolved.
Outdated

def _validate_params(self) -> None:
"""Validate initialization parameters.

Raises
------
InvalidSiteError
If no state is provided.
DateRangeError
If dates are invalid.
"""
if not self.state:
raise InvalidSiteError("state cannot be empty")

if self.start_date > self.end_date:
raise DateRangeError(
f"start_date ({self.start_date}) must be <= end_date ({self.end_date})"
)

def __init__(
self,
start_date: str,
end_date: str,
state: str = "ASOS",
variable: str | list[str] | None = None,
transformations: dict[str, Any] | None = None,
) -> None:
"""Fetch data from IEM ASOS API.

Parameters
----------
start_date : str
Start date in ISO 8601 format (YYYY-MM-DD).
end_date : str
End date in ISO 8601 format (YYYY-MM-DD).
state : str, optional
State abbreviation (e.g., 'TX' or 'ASOS' for all), by default "ASOS".
variable : str | list[str] | None, optional
Variable to fetch (e.g., 'p01i', 'tmpf', or ['p01i', 'tmpf']).
If None, fetches 'all'.
transformations : dict[str, Any], optional
Optional transformations to apply to the data.
"""
super().__init__()
self.state = state
self.network = "ASOS" if state.upper() == "ASOS" else f"{state.upper()}_ASOS"
Comment thread
sray014 marked this conversation as resolved.
Outdated

try:
self.start_date = pd.to_datetime(start_date)
self.end_date = pd.to_datetime(end_date)
except Exception as e:
raise DateRangeError(f"Could not parse date: {e}") from e

if variable is None:
Comment thread
sray014 marked this conversation as resolved.
self.data_vars = ["all"]
elif isinstance(variable, str):
self.data_vars = [variable]
else:
self.data_vars = variable

self.transformations = transformations

self._validate_params()
logger.debug(
f"Initialized {self.__class__.__name__}: "
f"network={self.network}, dates={self.start_date} to {self.end_date}, "
f"data={self.data_vars}"
)

def _build_url(self) -> str:
"""Build the IEM ASOS request URL with query parameters."""
params = [
("network", self.network),
("year1", self.start_date.year),
("month1", self.start_date.month),
("day1", self.start_date.day),
("year2", self.end_date.year),
("month2", self.end_date.month),
("day2", self.end_date.day),
("format", "comma"),
("latlon", "yes"),
]
for var in self.data_vars:
params.append(("data", var))
return f"{BASE_URL}?{urlencode(params)}"
Comment thread
sray014 marked this conversation as resolved.
Outdated

def _fetch(self, url: str) -> str:
"""Fetch CSV text from IEM via tiny_retriever."""
with wrap_errors(APIError, f"Failed to fetch ASOS data for {self.network}"):
return tiny_retriever.fetch(url, "text", timeout=self.timeout)

def _parse_response(self, text: str) -> pd.DataFrame:
"""Parse IEM CSV text into a DataFrame.

The first 5 rows of IEM ASOS output are comments (skiprows=5).
"""
df = pd.read_csv(StringIO(text), skiprows=5)
if df.empty:
logger.warning(f"ASOS returned no data for network {self.network}")
return pd.DataFrame()
Comment thread
sray014 marked this conversation as resolved.
Outdated
logger.debug(f"Fetched {len(df)} records from ASOS for {self.network}")
return df

def _reap(self) -> pd.DataFrame:
"""Fetch data from ASOS and return as a pandas DataFrame."""
logger.info(
f"Reaping ASOS data: network={self.network}, "
f"data={self.data_vars}"
)

url = self._build_url()
text = self._fetch(url)
df = self._parse_response(text)

if self.transformations and not df.empty:
df = apply_ts_transformations(df, self.transformations)

logger.info(f"Reaped {len(df)} records from {self.network}")
return df
180 changes: 180 additions & 0 deletions tests/test_reaping/test_asos_reapers.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
"""Tests for ASOS reapers."""

from __future__ import annotations

from unittest.mock import patch

import pandas as pd
import pytest

from cosecha.exceptions import APIError, DateRangeError, InvalidSiteError
from cosecha.reaping.asos import ASOSReaper


class TestASOSReaper:
"""Tests for ASOSReaper."""

def test_initialization_valid(self):
"""Test valid initialization."""
reaper = ASOSReaper(
state="TX",
variable="p01i",
start_date="2026-04-12",
end_date="2026-04-13",
)
assert reaper.state == "TX"
assert reaper.network == "TX_ASOS"
assert reaper.start_date == pd.Timestamp("2026-04-12")
assert reaper.end_date == pd.Timestamp("2026-04-13")
assert reaper.data_vars == ["p01i"]

def test_initialization_multiple_data_vars(self):
"""Test initialization with multiple data variables."""
reaper = ASOSReaper(
state="IA",
variable=["p01i", "tmpf"],
start_date="2026-04-12",
end_date="2026-04-13",
)
assert reaper.data_vars == ["p01i", "tmpf"]
assert reaper.network == "IA_ASOS"

def test_initialization_all_network(self):
"""Test initialization with ASOS state."""
reaper = ASOSReaper(
state="ASOS",
start_date="2026-04-12",
end_date="2026-04-13",
)
assert reaper.network == "ASOS"
assert reaper.data_vars == ["all"]

def test_empty_state(self):
"""Test initialization fails with empty state."""
with pytest.raises(InvalidSiteError, match="state cannot be empty"):
ASOSReaper(
state="",
start_date="2026-04-12",
end_date="2026-04-13",
)

def test_invalid_date_range(self):
"""Test initialization fails with invalid date range."""
with pytest.raises(DateRangeError):
ASOSReaper(
state="TX",
variable="p01i",
start_date="2026-04-13",
end_date="2026-04-12",
)

def test_invalid_start_date(self):
"""Test initialization fails with unparsable start date."""
with pytest.raises(DateRangeError, match="Could not parse date"):
ASOSReaper(
state="TX",
variable="p01i",
start_date="not-a-date",
end_date="2026-04-13",
)

@pytest.mark.network
def test_reap_network(self):
"""Test live network retrieval from ASOS."""
reaper = ASOSReaper(
state="TX",
variable="p01i",
start_date="2022-01-01",
end_date="2022-01-02",
)
harvested = reaper.reap()
assert len(harvested) > 0
assert "p01i" in harvested.columns
assert isinstance(harvested, pd.DataFrame)

@patch("cosecha.reaping.asos.tiny_retriever.fetch")
def test_reap_success(self, mock_fetch):
"""Test successful data retrieval."""
mock_fetch.return_value = (
"# comment 1\n"
"# comment 2\n"
"# comment 3\n"
"# comment 4\n"
"# comment 5\n"
"station,valid,lon,lat,p01i\n"
"AUS,2026-04-12 00:00,-97.6698,30.1945,0.01\n"
"AUS,2026-04-12 01:00,-97.6698,30.1945,0.02\n"
)

reaper = ASOSReaper(
state="TX",
variable="p01i",
start_date="2026-04-12",
end_date="2026-04-12",
)

harvested = reaper.reap()

assert len(harvested) == 2
assert isinstance(harvested, pd.DataFrame)
assert "p01i" in harvested.columns
assert "station" in harvested.columns
mock_fetch.assert_called_once()

url = mock_fetch.call_args[0][0]
assert "network=TX_ASOS" in url
assert "data=p01i" in url

@patch("cosecha.reaping.asos.tiny_retriever.fetch")
def test_reap_api_error(self, mock_fetch):
"""Test error handling when API fails."""
mock_fetch.side_effect = Exception("Connection failed")

reaper = ASOSReaper(
state="TX",
start_date="2026-04-12",
end_date="2026-04-12",
)

with pytest.raises(APIError, match="Failed to fetch ASOS data for TX_ASOS"):
reaper.reap()

mock_fetch.assert_called_once()

@patch("cosecha.reaping.asos.tiny_retriever.fetch")
def test_reap_empty_response(self, mock_fetch):
"""Test reap handles empty output after skipping comments."""
mock_fetch.return_value = (
"# 1\n# 2\n# 3\n# 4\n# 5\n"
"station,valid,lon,lat,p01i\n"
)

reaper = ASOSReaper(
state="TX",
start_date="2026-04-12",
end_date="2026-04-12",
)

harvested = reaper.reap()
assert isinstance(harvested, pd.DataFrame)
assert harvested.empty

@patch("cosecha.reaping.asos.tiny_retriever.fetch")
def test_reap_with_transformations(self, mock_fetch):
"""Test reap applies format transformations when not empty."""
mock_fetch.return_value = (
"# 1\n# 2\n# 3\n# 4\n# 5\n"
"station,p01i\n"
"AUS,0.01\n"
)

reaper = ASOSReaper(
state="TX",
start_date="2026-04-12",
end_date="2026-04-12",
transformations={"rename_columns": {"p01i": "precip"}},
)

harvested = reaper.reap()
assert "precip" in harvested.columns
assert "p01i" not in harvested.columns