diff --git a/modules/commands/rain_command.py b/modules/commands/rain_command.py index 1eaf1aa..1dc5e30 100644 --- a/modules/commands/rain_command.py +++ b/modules/commands/rain_command.py @@ -8,7 +8,7 @@ import asyncio import re import time from dataclasses import dataclass -from datetime import datetime, timedelta +from datetime import datetime, timedelta, timezone from typing import Any, Optional import requests @@ -408,6 +408,156 @@ def fetch_precip_series( return series +# --- NWS gridpoint precip source --------------------------------------------- +# WHY THIS EXISTS: the Open-Meteo *forecast model* (fetch_precip_series, above) +# smooths away scattered, pop-up convection, so the nowcast can miss rain that is +# actually happening. Observed near Nashville (36.16, -86.78): Open-Meteo reported +# 0.00 in / ~12% precip across the next 3 h while NWS's own gridpoint showed +# 65-74% probability with measurable QPF — and thunderstorms were occurring. The +# model-based push therefore never fired. NWS's gridpoint forecast is +# forecaster-adjusted and does capture convective chances, so for US points we +# prefer it (fetch_precip_series_nws) and fall back to Open-Meteo only where NWS +# has no coverage (outside the US) or the request fails. + +# NWS gridpoint "weather" type -> a representative WMO code, so precip_bucket_for_code() +# classifies the NWS series exactly like it classifies the Open-Meteo one. +_NWS_WEATHER_CODE = [ + ("thunderstorm", 95), + ("snow", 73), ("blowing_snow", 73), ("snow_showers", 73), + ("ice", 66), ("sleet", 66), ("freezing", 66), ("ice_pellets", 66), + ("drizzle", 53), + ("rain_showers", 81), ("showers", 81), + ("rain", 63), +] + + +def _iso_duration_hours(dur: str) -> int: + """Hours spanned by an ISO-8601 duration like 'PT6H', 'PT1H', 'P1DT6H' (min 1).""" + m = re.match(r"P(?:(\d+)D)?(?:T(?:(\d+)H)?(?:(\d+)M)?)?", dur or "") + if not m: + return 1 + days, hours, mins = (int(g) if g else 0 for g in m.groups()) + return max(1, days * 24 + hours + (1 if mins else 0)) + + +def _nws_hourly(values: Optional[list], *, divide: bool) -> dict: + """Map hour-start (naive UTC datetime) -> value from an NWS gridpoint property. + + NWS reports each property as time-bucketed values whose validTime is an ISO + interval like '2026-06-08T12:00:00+00:00/PT6H'. ``divide`` splits an + accumulation (e.g. 6-hour QPF) evenly across its hours; otherwise the period's + value is repeated for each hour (hourly PoP, the weather-type list). + """ + out: dict = {} + for v in values or []: + try: + start_s, _, dur = (v.get("validTime") or "").partition("/") + start = datetime.fromisoformat(start_s).astimezone(timezone.utc).replace(tzinfo=None) + except (TypeError, ValueError, AttributeError): + continue + n = _iso_duration_hours(dur) + raw = v.get("value") + share = (raw / n) if (divide and raw is not None) else raw + for k in range(n): + out[start + timedelta(hours=k)] = share + return out + + +def _nws_weather_code(value: Any) -> Optional[int]: + """Pick a representative WMO code from an NWS gridpoint ``weather`` value (list of segments).""" + if not value: + return None + blob = " ".join( + str(seg.get("weather") or "") for seg in value if isinstance(seg, dict) + ).lower() + if not blob.strip(): + return None + for needle, code in _NWS_WEATHER_CODE: + if needle in blob or needle.replace("_", " ") in blob: + return code + return 63 # precip of unknown type -> rain + + +def fetch_precip_series_nws( + session: Any, + lat: float, + lon: float, + *, + timeout: int = 10, + logger: Any = None, + pop_floor: int = 50, +) -> Optional[dict]: + """Build a precip nowcast series from the NWS gridpoint forecast (US only). + + Returns the same shape as fetch_precip_series (times/precip/codes/now/ + current_precip/current_code/step), or None when NWS has no coverage (e.g. + outside the US) so the caller can fall back to Open-Meteo. + + NWS exposes 6-hour QPF (mm) and hourly PoP (%). We build an hourly series in + which each hour's precip is its QPF share, but zeroed when that hour's PoP is + below ``pop_floor`` -- so the predicted rain-start tracks the hourly + probability rather than snapping to coarse 6-hour QPF boundaries, and a trace + of QPF at a low chance is not reported as rain. Times are naive UTC ISO strings + (they only need to be self-consistent: the nowcast works on relative minutes). + """ + headers = {"User-Agent": "(meshcore-bot, weather-nowcast)", "Accept": "application/geo+json"} + try: + pts = session.get( + f"https://api.weather.gov/points/{round(lat, 4)},{round(lon, 4)}", + headers=headers, timeout=timeout, + ) + if not pts.ok: + return None # no NWS coverage (outside the US) -> caller falls back to Open-Meteo + grid_url = (pts.json().get("properties") or {}).get("forecastGridData") + if not grid_url: + return None + gp = session.get(grid_url, headers=headers, timeout=timeout) + if not gp.ok: + return None + props = gp.json().get("properties") or {} + except (requests.exceptions.Timeout, requests.exceptions.ConnectionError) as e: + if logger: + logger.debug(f"NWS nowcast timeout/connection error: {e}") + return None + except (ValueError, KeyError, TypeError) as e: + if logger: + logger.debug(f"NWS nowcast parse error: {e}") + return None + + qpf = _nws_hourly((props.get("quantitativePrecipitation") or {}).get("values"), divide=True) + pop = _nws_hourly((props.get("probabilityOfPrecipitation") or {}).get("values"), divide=False) + wx = _nws_hourly((props.get("weather") or {}).get("values"), divide=False) + if not qpf and not pop: + return None + + now = datetime.now(timezone.utc).replace(tzinfo=None) + base = now.replace(minute=0, second=0, microsecond=0) + hours = [base + timedelta(hours=i) for i in range(0, 6)] # current hour + 5 ahead (covers the window) + + times: list[str] = [] + precip: list[Optional[float]] = [] + codes: list[Optional[int]] = [] + for h in hours: + p = pop.get(h) + q = qpf.get(h) + # Count an hour as precipitating only when NWS gives a real chance; the + # amount is its QPF share. (QPF is 6-hourly, PoP hourly -- PoP sets timing.) + amt = q if (q is not None and p is not None and p >= pop_floor) else 0.0 + times.append(h.isoformat(timespec="minutes")) + precip.append(amt) + codes.append(_nws_weather_code(wx.get(h)) if amt else None) + + return { + "times": times, + "precip": precip, + "codes": codes, + "now": now.isoformat(timespec="minutes"), + "current_precip": precip[0] if precip else None, + "current_code": codes[0] if codes else None, + "step": 60, + } + + @dataclass class NowcastResult: """Outcome of a precipitation nowcast analysis. @@ -892,9 +1042,20 @@ class RainCommand(BaseCommand): return (lat, lon, label, None) def _fetch_series(self, lat: float, lon: float) -> Optional[dict]: - """Fetch the precip series via the shared fetcher (own short-lived session).""" + """Fetch the precip series (own short-lived session). + + Prefers the NWS gridpoint (US) so a "!rain"/"!snow" matches the proactive + push and reflects the forecaster-adjusted convective chances the Open-Meteo + model can miss; falls back to Open-Meteo for non-US locations (no NWS + coverage) or on any failure. + """ session = self._create_retry_session() try: + series = fetch_precip_series_nws( + session, lat, lon, timeout=self.url_timeout, logger=self.logger, + ) + if series: + return series return fetch_precip_series( session, lat, lon, weather_model=self.weather_model, timeout=self.url_timeout, logger=self.logger, diff --git a/modules/service_plugins/weather_service.py b/modules/service_plugins/weather_service.py index d7bb0bc..a9d8b24 100644 --- a/modules/service_plugins/weather_service.py +++ b/modules/service_plugins/weather_service.py @@ -35,6 +35,7 @@ from ..commands.rain_command import ( decide_rain_notification, episode_probability_temp, fetch_precip_series, + fetch_precip_series_nws, format_amount_estimate, join_location, precip_descriptor, @@ -836,18 +837,33 @@ class WeatherService(BaseServicePlugin): """Fetch the precip nowcast for the bot's position and push if rain is incoming.""" try: loop = asyncio.get_event_loop() + # Prefer the NWS gridpoint (forecaster-adjusted QPF + PoP — it captures + # the convection the Open-Meteo model smooths away, which is why this + # push could stay silent during real rain). Fall back to Open-Meteo when + # NWS has no coverage (non-US) or the request fails. series = await loop.run_in_executor( None, - lambda: fetch_precip_series( + lambda: fetch_precip_series_nws( self.api_session, self.my_position_lat, self.my_position_lon, - weather_model=self.weather_model or "", timeout=10, logger=self.logger, cache_ttl=self.rain_nowcast_cache_seconds, ), ) + if not series: + series = await loop.run_in_executor( + None, + lambda: fetch_precip_series( + self.api_session, + self.my_position_lat, + self.my_position_lon, + weather_model=self.weather_model or "", + timeout=10, + logger=self.logger, + ), + ) if not series: return diff --git a/tests/unit/test_rain_nowcast.py b/tests/unit/test_rain_nowcast.py index 4f44f6f..216625a 100644 --- a/tests/unit/test_rain_nowcast.py +++ b/tests/unit/test_rain_nowcast.py @@ -10,6 +10,9 @@ from modules.commands.rain_command import ( RAIN_FAMILY, SNOW_FAMILY, NowcastResult, + _iso_duration_hours, # noqa: PLC2701 (testing internal helper) + _nws_hourly, # noqa: PLC2701 + _nws_weather_code, # noqa: PLC2701 _round5, # noqa: PLC2701 (testing internal helper) analyze_precip_nowcast, city_display_name, @@ -602,3 +605,43 @@ def test_region_note_fits_channel_budget(): worst_forecast = "🌧️ Heavy rain steady for 2h+ in Paris, France (est 0.5 in)" combined = f"{worst_forecast} {REGION_DEFAULT_NOTE}" assert len(combined.encode("utf-8")) <= budget + + +# --- NWS gridpoint source (fetch_precip_series_nws helpers) ------------------- + +def test_iso_duration_hours(): + assert _iso_duration_hours("PT1H") == 1 + assert _iso_duration_hours("PT6H") == 6 + assert _iso_duration_hours("P1DT6H") == 30 + assert _iso_duration_hours("PT3H") == 3 + assert _iso_duration_hours("") == 1 # unparseable -> at least 1 hour + assert _iso_duration_hours("PT30M") == 1 # sub-hour -> 1 + + +def test_nws_weather_code_classification(): + # NWS weather types map to WMO codes that classify into the same buckets. + assert precip_bucket_for_code(_nws_weather_code([{"weather": "thunderstorms"}])) == "thunder" + assert precip_bucket_for_code(_nws_weather_code([{"weather": "snow"}])) == "snow" + assert precip_bucket_for_code(_nws_weather_code([{"weather": "rain"}])) == "rain" + assert precip_bucket_for_code(_nws_weather_code([{"weather": "freezing_rain"}])) == "freezing" + # a thunderstorm-with-rain still classifies as thunder (priority order) + assert _nws_weather_code([{"weather": "rain"}, {"weather": "thunderstorms"}]) == 95 + assert _nws_weather_code([]) is None + assert _nws_weather_code([{"weather": None}]) is None + + +def test_nws_hourly_divide_and_repeat(): + # a 6-hour QPF accumulation is split evenly across its hours + spread = _nws_hourly( + [{"validTime": "2026-06-08T12:00:00+00:00/PT6H", "value": 6.0}], divide=True + ) + assert len(spread) == 6 + assert all(abs(v - 1.0) < 1e-9 for v in spread.values()) + # an hourly PoP is repeated (not divided) + pop = _nws_hourly( + [{"validTime": "2026-06-08T12:00:00+00:00/PT1H", "value": 70}], divide=False + ) + assert list(pop.values()) == [70] + # malformed / empty inputs are skipped, not fatal + assert _nws_hourly([{"value": 5}], divide=True) == {} + assert _nws_hourly(None, divide=False) == {} diff --git a/tests/unit/test_rain_proactive_e2e.py b/tests/unit/test_rain_proactive_e2e.py index 40ce27a..1e10db0 100644 --- a/tests/unit/test_rain_proactive_e2e.py +++ b/tests/unit/test_rain_proactive_e2e.py @@ -62,6 +62,12 @@ def build_service(series, monkeypatch, *, overrides=None): service._cached_rain_location = "Nashville, TN" # get_mesh_flood_scope lazily imports heavy deps; stub it. service.get_mesh_flood_scope = Mock(return_value=None) + # NWS gridpoint is tried first now; return None ("no coverage") so these + # source-agnostic nowcast-logic tests run on the canned Open-Meteo series. + monkeypatch.setattr( + "modules.service_plugins.weather_service.fetch_precip_series_nws", + lambda *a, **k: None, + ) monkeypatch.setattr( "modules.service_plugins.weather_service.fetch_precip_series", lambda *a, **k: series,