health check
This commit is contained in:
@@ -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,
|
||||
)
|
||||
|
||||
@@ -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
|
||||
|
||||
+27
-1
@@ -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()
|
||||
|
||||
+14
-1
@@ -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)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user