feat: public flood notifications over self-hosted ntfy
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
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
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).
This commit is contained in:
@@ -0,0 +1,132 @@
|
||||
"""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())
|
||||
Reference in New Issue
Block a user