diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index bdbfac7..41e7eaa 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -223,8 +223,13 @@ Identity (Universal Auth) и генерирует `.env` заново (`infisica ## 12. Наблюдаемость и эксплуатация -- **Мониторинг**: `antiplag_monitor.py` (cron 5 мин, 7 сервисов, email-алерт при смене - статуса). Опционально — Prometheus+Grafana (профиль `observability`, метрики API + Flower). +- **Мониторинг**: [`scripts/ops/antiplag_monitor.py`](../scripts/ops/antiplag_monitor.py) + (cron 5 мин на 1.32, email при смене статуса). Адреса берутся из прод-`.env`, а не + зашиты в код: прежняя версия месяц проверяла остановленный CT 108 и не видела + реального сервера эмбеддингов — постоянный ложный DOWN заглушал настоящие аварии. + Проверяются в том числе **воркеры Celery**: живой брокер ещё не значит работающую + систему (05.09 воркеры час простаивали при «зелёном» брокере). + Опционально — Prometheus+Grafana (профиль `observability`, метрики API + Flower). - **Бэкапы**: `pg_dump→gzip→MinIO`, cron 03:00, ротация 14. Проверка восстановления — `scripts/ops/pg_restore_verify.sh`. HA/DR — [DR-HA.md](DR-HA.md). - **Шкала загрузки источников** (админка → «Источники»): каждый запуск создаёт строку diff --git a/scripts/ops/antiplag_monitor.py b/scripts/ops/antiplag_monitor.py new file mode 100644 index 0000000..f0c0f20 --- /dev/null +++ b/scripts/ops/antiplag_monitor.py @@ -0,0 +1,204 @@ +#!/usr/bin/env python3 +"""Health-монитор anti-plagiarism: письмо при СМЕНЕ статуса сервиса, без спама. + +Запускается по cron на app-хосте (1.32) каждые 5 минут: + */5 * * * * /usr/bin/python3 /home/user/antiplag_monitor.py >> ~/.antiplag-monitor/run.log 2>&1 + +Живёт в репозитории намеренно: прошлая версия лежала только на сервере, ничем +не версионировалась и месяц проверяла топологию, которой уже нет — вечный +ложный DOWN на остановленном CT 108 и полная слепота к реальному серверу +эмбеддингов. Постоянный ложный сигнал хуже отсутствия сигнала: на его фоне +теряется настоящая авария. + +Что проверяется и почему именно это: +- API, PostgreSQL, Redis, RabbitMQ, MinIO, Elasticsearch — без них сервис не + принимает и не обрабатывает работы; +- Ollama на 1.40 — эмбеддинги для L3 (адрес брать из .env, а не хардкодить: + он уже дважды переезжал); +- воркеры Celery — самое важное дополнение: 05.09 брокер лежал час, воркеры + молчали, а прежний монитор рапортовал, что всё хорошо, потому что смотрел + только на TCP-порты. + +OpenRouter (L4) намеренно не проверяется: это внешний сервис за VPN-прокси, его +недоступность не ломает основную проверку — L4 лишь уточняет вердикт. +""" + +import json +import os +import smtplib +import socket +import ssl +import subprocess +import urllib.request +from datetime import datetime +from email.mime.text import MIMEText + +APP_DIR = os.environ.get("APP_DIR", "/home/user/anti-plagiarism") +STATE_FILE = os.path.expanduser("~/.antiplag-monitor/state.json") + +SMTP_HOST, SMTP_PORT = "mail.jze9mail.ru", 587 +SMTP_USER = os.environ.get("MONITOR_SMTP_USER", "noreply") +SMTP_PASS = os.environ.get("MONITOR_SMTP_PASS", "kGy3vgt6EuiUBMNDmV01Ng") +SMTP_FROM = "noreply@jze9mail.ru" +ALERT_TO = os.environ.get("MONITOR_ALERT_TO", "jze9programer@gmail.com") + + +def env_value(key: str, default: str = "") -> str: + """Значение из прод-.env (он генерируется из Infisical на каждом деплое). + + Адреса инфраструктуры читаем оттуда, а не хардкодим: сервер эмбеддингов уже + переезжал дважды, и монитор каждый раз оставался со старым адресом. + """ + path = os.path.join(APP_DIR, ".env") + try: + with open(path, encoding="utf-8") as fh: + for line in fh: + line = line.strip() + if line.startswith(f"{key}="): + return line.split("=", 1)[1].strip().strip("'\"") + except Exception: + pass + return default + + +def host_port_from_url(url: str, default_port: int) -> tuple[str, int]: + """Разобрать URL в пару (host, port), отбросив логин с паролем. + + В REDIS_URL и RABBITMQ_URL креды идут перед адресом + (`redis://:пароль@host:port/0`), и без их отсечения разбор падает. + """ + rest = url.split("//", 1)[-1].split("/", 1)[0] + rest = rest.rpartition("@")[2] or rest + if ":" in rest: + host, _, port = rest.partition(":") + return host, int(port or default_port) + return rest, default_port + + +def check_tcp(host: str, port: int) -> bool: + try: + with socket.create_connection((host, port), timeout=5): + return True + except Exception: + return False + + +def check_http(url: str) -> bool: + try: + with urllib.request.urlopen(url, timeout=8) as r: + return r.status == 200 + except Exception: + return False + + +def check_in_api_container(command: str) -> bool: + """Проверка изнутри контейнера — так же, как это видит приложение.""" + try: + out = subprocess.run( + ["docker", "exec", "antiplagiator-api", "sh", "-c", command], + capture_output=True, text=True, timeout=20, + ) + return out.stdout.strip() == "200" + except Exception: + return False + + +def check_celery_workers() -> bool: + """Отвечают ли воркеры на ping. + + Живой брокер ещё не значит работающую систему: воркер может не суметь к + нему подключиться и молча простаивать — именно так и было 05.09. + """ + try: + out = subprocess.run( + ["docker", "exec", "antiplagiator-api", "python", "-c", + "from app.core.celery_app import celery_app;" + "print(len(celery_app.control.inspect(timeout=5).ping() or {}))"], + capture_output=True, text=True, timeout=30, + ) + return int(out.stdout.strip() or 0) >= 4 # api ждёт 4 воркера + except Exception: + return False + + +def build_checks() -> list[tuple[str, callable]]: + """Список проверок; адреса — из прод-.env, чтобы не разъезжаться с реальностью.""" + pg_host = env_value("POSTGRES_HOST", "192.168.1.38") + redis_host, redis_port = host_port_from_url( + env_value("REDIS_URL", "redis://192.168.1.35:6379/0"), 6379) + minio_host, minio_port = host_port_from_url( + "http://" + env_value("MINIO_ENDPOINT", "192.168.1.21:9000"), 9000) + ollama_url = env_value("OLLAMA_URL", "http://192.168.1.40:11434") + + rabbit_host, rabbit_port = host_port_from_url( + env_value("RABBITMQ_URL", "amqp://guest:guest@192.168.20.82:5672/"), 5672) + + return [ + ("API", lambda: check_http("http://localhost:8000/health")), + ("PostgreSQL", lambda: check_tcp(pg_host, 5432)), + ("Redis", lambda: check_tcp(redis_host, redis_port)), + ("RabbitMQ", lambda: check_tcp(rabbit_host, rabbit_port)), + ("MinIO", lambda: check_tcp(minio_host, minio_port)), + ("Elasticsearch", lambda: check_in_api_container( + "curl -s -o /dev/null -w '%{http_code}' http://elasticsearch:9200")), + ("Embeddings", lambda: check_http(f"{ollama_url}/api/version")), + ("CeleryWorkers", check_celery_workers), + ] + + +def send_alert(subject: str, body: str) -> None: + msg = MIMEText(body, "plain", "utf-8") + msg["Subject"] = subject + msg["From"] = SMTP_FROM + msg["To"] = ALERT_TO + ctx = ssl.create_default_context() + with smtplib.SMTP(SMTP_HOST, SMTP_PORT, timeout=25) as s: + s.ehlo() + s.starttls(context=ctx) + s.ehlo() + s.login(SMTP_USER, SMTP_PASS) + s.send_message(msg) + + +def main() -> None: + results = {name: fn() for name, fn in build_checks()} + + prev: dict = {} + if os.path.exists(STATE_FILE): + try: + with open(STATE_FILE) as fh: + prev = json.load(fh) + except Exception: + prev = {} + + changes = [] + for name, ok in results.items(): + was = prev.get(name, True) # первый запуск считаем нормой + if was and not ok: + changes.append(f"УПАЛ: {name}") + elif not was and ok: + changes.append(f"восстановился: {name}") + + if changes: + ts = datetime.now().strftime("%Y-%m-%d %H:%M") + status = "\n".join(f" {'OK ' if v else 'DOWN'} {k}" for k, v in results.items()) + body = (f"Изменения статуса сервисов anti-plagiarism ({ts}):\n\n" + + "\n".join(changes) + "\n\nТекущий статус:\n" + status) + subject = f"[anti-plagiarism] {changes[0]}" + if len(changes) > 1: + subject += f" (+{len(changes) - 1})" + try: + send_alert(subject, body) + except Exception as e: + print(f"не удалось отправить алерт: {e}") + + os.makedirs(os.path.dirname(STATE_FILE), exist_ok=True) + with open(STATE_FILE, "w") as fh: + json.dump(results, fh) + + line = " ".join(f"{k}={'OK' if v else 'DOWN'}" for k, v in results.items()) + print(datetime.now().strftime("%H:%M"), "|", line) + + +if __name__ == "__main__": + main()