"""Notification backend abstraction. ntfy is the only implementation today, but call sites always go through get_notifier() / notify_flight_staged() so a second backend could be added later (a new class + registry entry) without touching the pipeline code that triggers notifications. """ import json import urllib.error import urllib.request from abc import ABC, abstractmethod import config import models import tokens class Notifier(ABC): @abstractmethod def send(self, title: str, body: str, url: str | None = None, actions: list[dict] | None = None) -> None: ... class NtfyNotifier(Notifier): def __init__(self, server_url: str, topic: str, auth_token: str | None = None): self.server_url = server_url.rstrip("/") self.topic = topic self.auth_token = auth_token def send(self, title: str, body: str, url: str | None = None, actions: list[dict] | None = None) -> None: if not self.topic: return # not configured yet - silently skip rather than crash the pipeline payload = {"topic": self.topic, "title": title, "message": body} if url: payload["click"] = url if actions: payload["actions"] = actions req = urllib.request.Request( self.server_url, data=json.dumps(payload).encode("utf-8"), headers={"Content-Type": "application/json"}, ) if self.auth_token: req.add_header("Authorization", f"Bearer {self.auth_token}") try: with urllib.request.urlopen(req, timeout=15): pass except (urllib.error.HTTPError, urllib.error.URLError): pass # best-effort; a failed push shouldn't crash the pipeline run def get_notifier() -> Notifier: settings = models.get_all_settings() # Only one backend exists today; settings.notify_backend is read here so # a second implementation can be switched in later without changing # call sites. return NtfyNotifier( server_url=settings["ntfy_server_url"], topic=settings["ntfy_topic"], auth_token=settings["ntfy_auth_token"] or None, ) def flight_url(flight_id: int) -> str | None: if not config.PUBLIC_BASE_URL: return None return f"{config.PUBLIC_BASE_URL}/flights/{flight_id}" def notify_flight_staged(row: dict) -> None: flight_id = row["id"] route = f"{row.get('from_iata') or '?'}→{row.get('to_iata') or '?'}" flight_number = f"{row.get('operating_carrier_iata') or ''}{row.get('flight_number') or ''}" title = f"{flight_number or 'Flight'} staged for review" if row["status"] == "track_unavailable": summary = "no track found" else: count = row.get("track_point_count") or 0 summary = f"track found, {count} points" body = f"{route} · {row.get('flight_date') or '?'} · {summary}" review_url = flight_url(flight_id) actions = [] if review_url: actions.append({"action": "view", "label": "Investigate", "url": review_url}) base = config.PUBLIC_BASE_URL if base: approve_token = tokens.sign_token(flight_id, "approve") reject_token = tokens.sign_token(flight_id, "reject") actions.append({ "action": "http", "label": "Approve", "method": "POST", "clear": True, "url": f"{base}/flights/{flight_id}/approve?token={approve_token}", }) actions.append({ "action": "http", "label": "Deny", "method": "POST", "clear": True, "url": f"{base}/flights/{flight_id}/reject?token={reject_token}", }) get_notifier().send(title=title, body=body, url=review_url, actions=actions or None)