"""Finding, downloading, cropping, and staging meteorology files."""
from __future__ import annotations
import logging
import shutil
from pathlib import Path
from typing import TYPE_CHECKING, cast
import pandas as pd
from pandas.tseries.frequencies import to_offset
from stilt.config import MetConfig
from stilt.config.meteorology import arlmet_sources
from stilt.errors import MeteorologyError
if TYPE_CHECKING:
from arlmet.sources import MeteorologySource as ArlmetSource
logger = logging.getLogger(__name__)
[docs]
class MetStream:
"""
Meteorology files for one met stream, found locally or downloaded.
Without ``source``, files are found under ``directory`` from
``file_format`` and ``file_tres``. With ``source`` set to an arlmet
source name, arlmet downloads the files, cropping them as it goes when
subgridding is on. Local files are cropped with
``arlmet.extract_subset`` into ``subgrid_dir``, which all simulations
share.
Parameters
----------
name : str
Name of the met in the config.
config : MetConfig
Its settings.
"""
def __init__(self, name: str, config: MetConfig):
self.name = name
self.config = config
#: ``config.directory``, made absolute.
self.directory = config.directory.expanduser().resolve()
self._arlmet_source: ArlmetSource | None = None
# ------------------------------------------------------------------
# Internal helpers
# ------------------------------------------------------------------
@staticmethod
def _dedupe_matched_files(paths: list[Path]) -> list[Path]:
"""Resolve symlinks and drop duplicate files, sorted by name."""
return sorted(
dict.fromkeys(path.resolve() for path in paths),
key=lambda path: path.name,
)
def _get_arlmet_source(self) -> ArlmetSource:
"""Return the arlmet source, building it on first use."""
if self._arlmet_source is None:
assert self.config.source is not None
cls = arlmet_sources()[self.config.source]
self._arlmet_source = cls(**self.config.source_kwargs)
return self._arlmet_source
def _effective_bbox(self) -> tuple[float, float, float, float]:
"""Return ``(west, south, east, north)`` of the subgrid bounds plus the buffer."""
b = self.config.subgrid_bounds
buf = self.config.subgrid_buffer
if b is None:
raise ValueError("subgrid_bounds is required to compute effective bbox.")
if buf is None or buf < 0:
raise ValueError("subgrid_buffer must be a non-negative number.")
return (b.xmin - buf, b.ymin - buf, b.xmax + buf, b.ymax + buf)
def _resolved_subgrid_dir(self) -> Path:
"""Return the directory for cropped files, ``<directory>/subgrid`` by default."""
if self.config.subgrid_dir is None:
return self.directory / "subgrid"
return self.config.subgrid_dir.expanduser().resolve()
def _level_indices(self) -> list[int] | None:
"""Return the indices of the lowest ``subgrid_levels`` levels, or None to keep all."""
if self.config.subgrid_levels is None:
return None
return list(range(self.config.subgrid_levels))
# ------------------------------------------------------------------
# File resolution
# ------------------------------------------------------------------
def _fetch_from_source(self, r_time: pd.Timestamp, n_hours: int) -> list[Path]:
"""Return the files for a run from the arlmet source, downloading any not yet in ``directory``."""
sim_end = r_time + pd.Timedelta(hours=n_hours)
t_start: pd.Timestamp = min(r_time, sim_end) # type: ignore[assignment]
t_end: pd.Timestamp = max(r_time, sim_end) # type: ignore[assignment]
bbox = self._effective_bbox() if self.config.subgrid_enable else None
source = self._get_arlmet_source()
try:
files = source.fetch(
t_start,
t_end,
local_dir=self.directory,
backend=self.config.backend,
bbox=bbox,
)
except ImportError as exc:
raise ImportError(
f"{exc}\n\n"
"Downloading meteorology requires the cloud extra. "
"Install with: pip install pystilt[cloud]"
) from exc
n_files = len(files)
if n_files == 0 or n_files < self.config.n_min:
raise MeteorologyError(
f"Insufficient number of meteorological files found. "
f"Found: {n_files}, Required: {self.config.n_min}."
)
return files
[docs]
def required_files(self, r_time, n_hours: int) -> list[Path]:
"""
Return the met files that cover one simulation.
Parameters
----------
r_time : datetime-like
Receptor time.
n_hours : int
Simulation length in hours, negative for backward runs.
Returns
-------
list of Path
Raises
------
MeteorologyError
Fewer than ``n_min`` files were found.
"""
_r_time = cast(pd.Timestamp, pd.Timestamp(r_time))
if self.config.source is not None:
return self._fetch_from_source(_r_time, n_hours)
# Archive-glob mode
sim_end = _r_time + pd.Timedelta(hours=n_hours)
assert isinstance(sim_end, pd.Timestamp) # not NaT: _r_time is a time
earlier = min(_r_time, sim_end)
later = max(_r_time, sim_end)
file_format, file_tres = self.config.file_format, self.config.file_tres
# MetConfig requires both when there is no source.
assert file_format is not None and file_tres is not None
tres = to_offset(pd.to_timedelta(file_tres)).freqstr
met_start = earlier.floor(tres)
met_end = later
if n_hours < 0:
met_end_ceil = later.ceil(tres)
# As in STILT-R: a release in the last hour of a file interpolates
# against the next file's first hour; anywhere else it doesn't.
if later.floor("h") + pd.Timedelta(hours=1) == met_end_ceil: # type: ignore[arg-type]
met_end = met_end_ceil
met_times = pd.date_range(met_start, met_end, freq=tres)
patterns = list(dict.fromkeys(t.strftime(file_format) for t in met_times))
files: list[Path] = []
missing: list[str] = []
for pattern in patterns:
# Backup copies (name~<timestamp>~, name.~1~, name~) all end in "~".
matches = [
p
for p in self.directory.rglob(f"{pattern}*")
if p.is_file() and ".lock" not in p.name and not p.name.endswith("~")
]
if matches:
files.extend(matches)
else:
missing.append(pattern)
files = self._dedupe_matched_files(files)
n_files = len(files)
if n_files == 0 or n_files < self.config.n_min:
detail = ""
if missing:
examples = ", ".join(missing[:3])
detail = f" Patterns not found in {self.directory}: {examples}."
raise MeteorologyError(
f"Insufficient number of meteorological files found. "
f"Found: {n_files}, Required: {self.config.n_min}.{detail}"
)
if missing:
examples = ", ".join(missing[:3])
logger.warning(
"Met patterns not found (simulation may lack temporal coverage): %s",
examples,
)
return files
[docs]
def stage_files_for_simulation(
self,
*,
r_time,
n_hours: int,
target_dir: Path | str,
) -> list[Path]:
"""Find the met files for one simulation and link them into ``target_dir``."""
return self._stage_files(
self.required_files(r_time=r_time, n_hours=n_hours),
target_dir=target_dir,
)
def _stage_files(self, files: list[Path], target_dir: Path | str) -> list[Path]:
"""
Link met files into ``target_dir``, copying when a link fails.
With subgridding on and no ``source``, each file is cropped into
``subgrid_dir`` first and the cropped copy is linked. Downloaded files
were already cropped.
"""
# Resolve subgridded paths for archive-mode subsetting
if self.config.subgrid_enable and self.config.source is None:
files = self._subset_archive_files(files)
target = Path(target_dir)
target.mkdir(parents=True, exist_ok=True)
staged: list[Path] = []
staged_sources: dict[Path, Path] = {}
for src in files:
src = Path(src)
resolved_src = src.resolve()
if src.parent == target:
if src not in staged_sources:
staged_sources[src] = resolved_src
staged.append(src)
continue
dst = target / src.name
existing = staged_sources.get(dst)
if existing is not None:
if existing != resolved_src:
logger.warning(
"met source has duplicate basename %s at %s and %s; staging %s",
src.name,
existing,
resolved_src,
existing,
)
continue
staged_sources[dst] = resolved_src
if dst.exists() or dst.is_symlink():
staged.append(dst)
continue
try:
dst.symlink_to(resolved_src)
except OSError:
shutil.copy2(src, dst)
staged.append(dst)
return staged
def _subset_archive_files(self, files: list[Path]) -> list[Path]:
"""Crop local files into ``subgrid_dir``, reusing crops that already exist."""
from arlmet import extract_subset
subgrid_dir = self._resolved_subgrid_dir()
subgrid_dir.mkdir(parents=True, exist_ok=True)
bbox = self._effective_bbox()
levels = self._level_indices()
subsetted: list[Path] = []
for src in files:
cache_path = subgrid_dir / src.name
if not cache_path.exists():
logger.info("Subsetting %s → %s", src.name, cache_path)
extract_subset(src, cache_path, bbox=bbox, levels=levels)
subsetted.append(cache_path)
return subsetted