Files
grabowski 777b230baf
CI / Format & lint (push) Successful in 11s
Security / Static analysis (push) Successful in 12s
CI / Test suite (push) Successful in 18s
Docs / Validate documentation (push) Successful in 11s
Security / Dependency vulnerabilities (push) Successful in 1m30s
Security / License report (push) Successful in 47s
feat: public flood notifications over self-hosted ntfy
Anyone can now get push alerts on their phone without an account: the
monitor publishes to an ntfy server (one Go binary, ~30 MB RSS) and
subscribers pick topics in the free iOS/Android/web app.

Semantics are transitions, never state. One message when a gauge crosses
its warning or danger threshold, one all-clear when it drops back (0.10 m
hysteresis), nothing while it sits above. A three-day flood is two
messages; a quiet season is zero. Topics: ping-warning / ping-danger
(basin digest), ping-<station>-warning / -danger, ping-p1-outlook (opt-in:
model P(warning within 24 h) at P.1 rises through 50 %, clears below 25 %,
message says it is experimental), ping-status (feed stale >= 3 h /
recovered). Priority 5 on danger so it rings through Do Not Disturb.

src/notify.py runs once per collection cycle in the API process (leader
only, after the forecast precompute, same data the dashboard shows).
Last-sent state lives in a notification_state table so a restart never
re-sends; a failed publish leaves state untouched so the crossing is
retried next cycle instead of lost. Off unless NTFY_SERVER is set.

Dashboard: a "Get alerts" button (only when configured) opens a panel
with the server, per-topic cards, ntfy:// deep links and web links, app
store links and a disclaimer. EN + TH. GET /api/notifications feeds it.

scripts/install_ntfy.sh: .deb install, server.yml (loopback listen,
anonymous read, token-only write scoped to ping-*, 72 h cache, signup/
login/metrics off, tight visitor limits), systemd, user + token, .env.
Verified against ntfy 2.28.0: anon publish 403, token publish 200, token
on foreign topic 403, anon read 200, and a seeded crossing through the
real _notify_transitions path arrived in the topic with priority, tags,
click and action button. docs/NOTIFICATIONS.md has the deployment and
reverse-proxy notes. Tests: 10 for the state machine (159 total).
2026-09-12 00:18:38 +02:00

133 lines
4.2 KiB
Python

"""Drive the production notify path in-process: startup init -> seeded readings
-> forecast cache -> _notify_transitions -> sqlite state -> real ntfy."""
import asyncio
import datetime
import json
import os
import sys
import requests
os.environ.update(
DB_TYPE="sqlite",
WATER_DB_PATH=os.path.join(os.environ["LOCALAPPDATA"], "Temp", "smoke3.db"),
NTFY_SERVER="http://127.0.0.1:2586",
NTFY_TOKEN=os.environ.get("NTFY_TOKEN", ""),
NTFY_TOPIC_PREFIX="ping",
)
for f in ("smoke3.db",):
p = os.path.join(os.environ["LOCALAPPDATA"], "Temp", f)
if os.path.exists(p):
os.remove(p)
from src import web_api # noqa: E402
from src.config import Config # noqa: E402
assert Config.NTFY_SERVER
async def main():
# what the lifespan does at startup, minus the scheduler
from src import notify as notify_mod
from src.forecast_history import ForecastHistoryStore
from src.water_scraper_v3 import EnhancedWaterMonitorScraper
db_config = Config.get_database_config()
web_api.app_state["scraper"] = EnhancedWaterMonitorScraper(db_config)
store = ForecastHistoryStore(db_config["connection_string"], db_config["type"])
store.connect()
web_api.app_state["forecast_store"] = store
state = notify_mod.NotificationState(store.engine, store.db_type)
pub = notify_mod.NtfyPublisher(
Config.NTFY_SERVER, prefix=Config.NTFY_TOPIC_PREFIX, token=Config.NTFY_TOKEN
)
web_api.app_state["notify"] = (pub, state)
scraper = web_api.app_state["scraper"]
now = datetime.datetime.now().replace(minute=0, second=0, microsecond=0)
def seed(level_p1, level_p103, ts):
rows = [
{
"station_code": "P.1",
"station_id": 1,
"timestamp": ts,
"water_level": level_p1,
"discharge": 400.0,
"station_name_en": "Nawarat Bridge",
"station_name_th": "สะพานนวรัฐ",
"discharge_percent": 30.0,
"status": "active",
},
{
"station_code": "P.103",
"station_id": 2,
"timestamp": ts,
"water_level": level_p103,
"discharge": 300.0,
"station_name_en": "Ring Road 3",
"station_name_th": "วงแหวน 3",
"discharge_percent": 20.0,
"status": "active",
},
]
scraper.db_adapter.save_measurements(rows)
def forecast(p):
with web_api.FORECAST_CACHE_LOCK:
web_api.FORECAST_CACHE["all"] = (
0,
[
{
"station_code": "P.1",
"horizon_hours": 24,
"p_warning": p,
"predicted_max_level": 3.9,
"source": "model",
}
],
)
def poll(topic):
out = []
for line in (
requests.get(f"{Config.NTFY_SERVER}/{topic}/json?poll=1", timeout=5)
.text.strip()
.splitlines()
):
m = json.loads(line)
if m.get("event") == "message":
out.append(m.get("title") or m.get("message", "")[:40])
return out
# cycle 1: quiet
seed(1.6, 3.2, now - datetime.timedelta(hours=2))
forecast(0.02)
await web_api._notify_transitions()
# cycle 2: P.1 crosses warning, model outlook on
seed(3.75, 3.3, now - datetime.timedelta(hours=1))
forecast(0.7)
await web_api._notify_transitions()
# cycle 3: same state -> silence
seed(3.80, 3.3, now)
forecast(0.65)
await web_api._notify_transitions()
print("ping-p1-warning:", poll("ping-p1-warning"))
print("ping-warning: ", poll("ping-warning"))
print("ping-p1-outlook:", poll("ping-p1-outlook"))
print("ping-p103-warning:", poll("ping-p103-warning"))
from sqlalchemy import text
with store.engine.connect() as c:
print(
"state table:",
c.execute(
text("SELECT key, state, value FROM notification_state ORDER BY key")
).fetchall(),
)
asyncio.run(main())