diff --git a/docker-compose.yaml b/docker-compose.yaml index 426645c..7b3024c 100644 --- a/docker-compose.yaml +++ b/docker-compose.yaml @@ -11,6 +11,12 @@ services: - "/etc/hosts:/etc/hosts" ports: - 5004:5000 + healthcheck: + test: ["CMD", "python", "-c", "import urllib.request,sys; import json; r=urllib.request.urlopen('http://127.0.0.1:5000/healthz', timeout=3); sys.exit(0 if r.status==200 else 1)"] + interval: 30s + timeout: 5s + retries: 3 + start_period: 30s logging: driver: json-file options: diff --git a/src/health.py b/src/health.py new file mode 100644 index 0000000..aa25280 --- /dev/null +++ b/src/health.py @@ -0,0 +1,91 @@ +import threading +import time +from dataclasses import dataclass + + +@dataclass +class HealthSnapshot: + now: float + uptime_sec: float + telegram_last_ok_ago_sec: float | None + telegram_last_error_ago_sec: float | None + max_last_ok_ago_sec: float | None + max_last_error_ago_sec: float | None + max_last_event_ago_sec: float | None + + telegram_healthy: bool + max_healthy: bool + overall_healthy: bool + + +class HealthState: + def __init__(self, *, unhealthy_after_sec: float = 15 * 60) -> None: + self._lock = threading.Lock() + self._started_at = time.time() + self._unhealthy_after_sec = float(unhealthy_after_sec) + + self._telegram_last_ok: float | None = None + self._telegram_last_error: float | None = None + + self._max_last_ok: float | None = None + self._max_last_error: float | None = None + self._max_last_event: float | None = None + + def mark_telegram_ok(self) -> None: + with self._lock: + self._telegram_last_ok = time.time() + + def mark_telegram_error(self) -> None: + with self._lock: + self._telegram_last_error = time.time() + + def mark_max_ok(self) -> None: + with self._lock: + self._max_last_ok = time.time() + + def mark_max_error(self) -> None: + with self._lock: + self._max_last_error = time.time() + + def mark_max_event(self) -> None: + with self._lock: + self._max_last_event = time.time() + + def snapshot(self) -> HealthSnapshot: + now = time.time() + with self._lock: + started_at = self._started_at + unhealthy_after = self._unhealthy_after_sec + + t_ok = self._telegram_last_ok + t_err = self._telegram_last_error + m_ok = self._max_last_ok + m_err = self._max_last_error + m_evt = self._max_last_event + + def ago(ts: float | None) -> float | None: + if ts is None: + return None + return max(0.0, now - ts) + + telegram_last_ok_ago = ago(t_ok) + max_last_ok_ago = ago(m_ok) + + telegram_healthy = telegram_last_ok_ago is not None and telegram_last_ok_ago <= unhealthy_after + max_healthy = max_last_ok_ago is not None and max_last_ok_ago <= unhealthy_after + + overall_healthy = telegram_healthy and max_healthy + + return HealthSnapshot( + now=now, + uptime_sec=max(0.0, now - started_at), + telegram_last_ok_ago_sec=telegram_last_ok_ago, + telegram_last_error_ago_sec=ago(t_err), + max_last_ok_ago_sec=max_last_ok_ago, + max_last_error_ago_sec=ago(m_err), + max_last_event_ago_sec=ago(m_evt), + telegram_healthy=telegram_healthy, + max_healthy=max_healthy, + overall_healthy=overall_healthy, + ) + diff --git a/src/health_web.py b/src/health_web.py new file mode 100644 index 0000000..ac0f1a3 --- /dev/null +++ b/src/health_web.py @@ -0,0 +1,79 @@ +import json +import logging +import threading +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from typing import Any + +from health import HealthState + +logger = logging.getLogger(__name__) + + +class _Handler(BaseHTTPRequestHandler): + health: HealthState + + def log_message(self, format: str, *args: Any) -> None: # noqa: A003 + # Убираем спам от http.server, оставляем только наши логи. + logger.debug("health_http " + format, *args) + + def _send_json(self, status: int, payload: dict[str, Any]) -> None: + body = json.dumps(payload, ensure_ascii=False).encode("utf-8") + self.send_response(status) + self.send_header("Content-Type", "application/json; charset=utf-8") + self.send_header("Content-Length", str(len(body))) + self.end_headers() + self.wfile.write(body) + + def do_GET(self) -> None: # noqa: N802 + if self.path in {"/livez", "/live", "/"}: + snap = self.health.snapshot() + self._send_json( + 200, + { + "status": "live", + "uptime_sec": snap.uptime_sec, + }, + ) + return + + if self.path in {"/healthz", "/health"}: + snap = self.health.snapshot() + status = 200 if snap.overall_healthy else 503 + self._send_json( + status, + { + "status": "ok" if snap.overall_healthy else "unhealthy", + "telegram": { + "healthy": snap.telegram_healthy, + "last_ok_ago_sec": snap.telegram_last_ok_ago_sec, + "last_error_ago_sec": snap.telegram_last_error_ago_sec, + }, + "max": { + "healthy": snap.max_healthy, + "last_ok_ago_sec": snap.max_last_ok_ago_sec, + "last_error_ago_sec": snap.max_last_error_ago_sec, + "last_event_ago_sec": snap.max_last_event_ago_sec, + }, + "uptime_sec": snap.uptime_sec, + }, + ) + return + + self._send_json(404, {"error": "not_found"}) + + +def start_health_server(*, host: str, port: int, health: HealthState) -> threading.Thread: + class Handler(_Handler): + pass + + Handler.health = health + server = ThreadingHTTPServer((host, port), Handler) + + def _run() -> None: + logger.info("Health server listening on %s:%s", host, port) + server.serve_forever(poll_interval=0.5) + + thread = threading.Thread(target=_run, name="health-http", daemon=True) + thread.start() + return thread + diff --git a/src/main.py b/src/main.py index 69a25b1..4483405 100644 --- a/src/main.py +++ b/src/main.py @@ -6,6 +6,8 @@ from dotenv import load_dotenv from bridge import MaxToTelegramBridge from config import load_settings +from health import HealthState +from health_web import start_health_server from reverse_bridge import TelegramToMaxBridge from storage import BridgeStorage from telegram_api import TelegramClient @@ -44,16 +46,22 @@ def main() -> None: ) storage = BridgeStorage(settings.sqlite_path) bridge = MaxToTelegramBridge(max_client=max_client, telegram=telegram_client, storage=storage) - reverse_bridge = TelegramToMaxBridge(max_client=max_client, telegram=telegram_client, storage=storage) + health = HealthState(unhealthy_after_sec=15 * 60) + start_health_server(host="0.0.0.0", port=5000, health=health) + + reverse_bridge = TelegramToMaxBridge(max_client=max_client, telegram=telegram_client, storage=storage, health=health) logger = logging.getLogger("max2telegram") @max_client.on_start async def on_start() -> None: logger.info("Max client started as %s", max_client.me.id) + health.mark_max_ok() asyncio.create_task(reverse_bridge.start()) + asyncio.create_task(_max_probe_loop(max_client=max_client, health=health)) @max_client.on_message() async def on_message(message: Message) -> None: + health.mark_max_event() try: await bridge.forward_message(message) except Exception: @@ -62,5 +70,23 @@ def main() -> None: asyncio.run(max_client.start()) +async def _max_probe_loop(*, max_client: MaxClient, health: HealthState) -> None: + # Best-effort контроль соединения: периодически дергаем API. + # Если PyMax разорвет соединение/сломается сессия, это обычно проявится как исключение. + await asyncio.sleep(2) + while True: + try: + me = getattr(max_client, "me", None) + my_id = getattr(me, "id", None) + if my_id is not None: + await max_client.get_user(user_id=my_id) + health.mark_max_ok() + except Exception: + health.mark_max_error() + logger = logging.getLogger("max2telegram") + logger.exception("Max probe failed") + await asyncio.sleep(60) + + if __name__ == "__main__": main() diff --git a/src/reverse_bridge.py b/src/reverse_bridge.py index 6ee4f91..e13791f 100644 --- a/src/reverse_bridge.py +++ b/src/reverse_bridge.py @@ -9,6 +9,7 @@ from pymax.files import Photo, Video from storage import BridgeStorage from telegram_api import TelegramClient +from health import HealthState logger = logging.getLogger(__name__) @@ -53,10 +54,18 @@ class _MediaGroupBuffer: class TelegramToMaxBridge: - def __init__(self, *, max_client: MaxClient, telegram: TelegramClient, storage: BridgeStorage) -> None: + def __init__( + self, + *, + max_client: MaxClient, + telegram: TelegramClient, + storage: BridgeStorage, + health: "HealthState | None" = None, + ) -> None: self._max_client = max_client self._telegram = telegram self._storage = storage + self._health = health self._max_title_to_id: dict[str, int] = {} self._bot_id: str | None = None @@ -77,8 +86,12 @@ class TelegramToMaxBridge: while True: try: updates = await self._telegram.get_updates(offset=self._offset, timeout=25, limit=100) + if self._health: + self._health.mark_telegram_ok() await self._handle_updates(updates) except Exception: + if self._health: + self._health.mark_telegram_error() logger.exception("Telegram polling loop error") await asyncio.sleep(2)