import asyncio import hashlib import hmac import json import os import random import sqlite3 from contextlib import asynccontextmanager from dataclasses import dataclass from datetime import date, datetime, timedelta, timezone from typing import Any import httpx from fastapi import FastAPI, Header, HTTPException, Request DB_PATH = os.getenv("BOT_DB_PATH", "botreviewer.sqlite3") TELEGRAM_BOT_TOKEN = os.getenv("TELEGRAM_BOT_TOKEN", "") GITEA_BASE_URL = os.getenv("GITEA_BASE_URL", "").rstrip("/") GITEA_TOKEN = os.getenv("GITEA_TOKEN", "") GITEA_WEBHOOK_SECRET = os.getenv("GITEA_WEBHOOK_SECRET", "") REMINDER_TIME = os.getenv("REMINDER_TIME", "09:00") # Формат ЧЧ:ММ AUTO_ASSIGNED_LABEL = os.getenv("AUTO_ASSIGNED_LABEL", "auto-assigned") try: AUTO_ASSIGN_REVIEWERS_COUNT = max(1, int(os.getenv("AUTO_ASSIGN_REVIEWERS_COUNT", "2"))) except ValueError: AUTO_ASSIGN_REVIEWERS_COUNT = 2 try: APPROVALS_REQUIRED_FOR_MERGE = max(1, int(os.getenv("APPROVALS_REQUIRED_FOR_MERGE", "2"))) except ValueError: APPROVALS_REQUIRED_FOR_MERGE = 2 ALLOW_SELF_ASSIGN = os.getenv("ALLOW_SELF_ASSIGN", "false").lower() == "true" TELEGRAM_WEBHOOK_SECRET = os.getenv("TELEGRAM_WEBHOOK_SECRET", "") TELEGRAM_ALLOWED_CHAT_IDS = { int(v.strip()) for v in os.getenv("TELEGRAM_ALLOWED_CHAT_IDS", "").split(",") if v.strip().isdigit() } TELEGRAM_ALLOWED_USERNAMES = { v.strip().lstrip("@").lower() for v in os.getenv("TELEGRAM_ALLOWED_USERNAMES", "").split(",") if v.strip() } GITEA_COUNT_REPOS = { v.strip() for v in os.getenv("GITEA_COUNT_REPOS", "").split(",") if v.strip() } AUTO_ASSIGN_EXCLUDED_AUTHORS = { v.strip().lower() for v in os.getenv("AUTO_ASSIGN_EXCLUDED_AUTHORS", "").split(",") if v.strip() } NOTIFY_ON_WEEKENDS = os.getenv("NOTIFY_ON_WEEKENDS", "true").lower() == "true" MSK_TZ = timezone(timedelta(hours=3)) @dataclass class RegisteredUser: telegram_chat_id: int telegram_username: str | None email: str gitea_login: str def utc_now() -> datetime: return datetime.now(timezone.utc) def msk_now() -> datetime: return datetime.now(MSK_TZ) def parse_date(value: str) -> date: return datetime.strptime(value, "%Y-%m-%d").date() def _log_ts() -> str: return msk_now().strftime("%Y-%m-%d %H:%M:%S MSK") def log_info(msg: str) -> None: print(f"[{_log_ts()}] [info] {msg}") def log_warn(msg: str) -> None: print(f"[{_log_ts()}] [warn] {msg}") def log_error(context: str, exc: Exception) -> None: # Не логируем str(exc), чтобы не утекали токены/URL/чувствительные детали. print(f"[{_log_ts()}] [error] {context}: {type(exc).__name__}") def is_user_allowed(chat_id: int, username: str | None) -> bool: if not TELEGRAM_ALLOWED_CHAT_IDS and not TELEGRAM_ALLOWED_USERNAMES: return True username_norm = (username or "").lstrip("@").lower() return chat_id in TELEGRAM_ALLOWED_CHAT_IDS or (username_norm and username_norm in TELEGRAM_ALLOWED_USERNAMES) class Storage: def __init__(self, db_path: str) -> None: self.db_path = db_path self._ensure_schema() def _connect(self) -> sqlite3.Connection: conn = sqlite3.connect(self.db_path) conn.row_factory = sqlite3.Row return conn def _ensure_schema(self) -> None: with self._connect() as conn: conn.executescript( """ CREATE TABLE IF NOT EXISTS users ( telegram_chat_id INTEGER PRIMARY KEY, telegram_username TEXT, email TEXT UNIQUE NOT NULL, gitea_login TEXT UNIQUE NOT NULL, created_at TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS away_periods ( id INTEGER PRIMARY KEY AUTOINCREMENT, telegram_chat_id INTEGER NOT NULL, date_from TEXT NOT NULL, date_to TEXT NOT NULL ); CREATE TABLE IF NOT EXISTS assignments ( id INTEGER PRIMARY KEY AUTOINCREMENT, repo_full_name TEXT NOT NULL, pr_number INTEGER NOT NULL, author_login TEXT NOT NULL, reviewer_login TEXT NOT NULL, status TEXT NOT NULL, created_at TEXT NOT NULL, closed_at TEXT ); CREATE UNIQUE INDEX IF NOT EXISTS ux_assignments_repo_pr_reviewer ON assignments(repo_full_name, pr_number, reviewer_login); CREATE TABLE IF NOT EXISTS kv ( key TEXT PRIMARY KEY, value TEXT NOT NULL ); """ ) def upsert_user(self, user: RegisteredUser) -> None: with self._connect() as conn: conn.execute( """ INSERT INTO users (telegram_chat_id, telegram_username, email, gitea_login, created_at) VALUES (?, ?, ?, ?, ?) ON CONFLICT(telegram_chat_id) DO UPDATE SET telegram_username=excluded.telegram_username, email=excluded.email, gitea_login=excluded.gitea_login """, ( user.telegram_chat_id, user.telegram_username, user.email.lower(), user.gitea_login, utc_now().isoformat(), ), ) def set_away(self, telegram_chat_id: int, date_from: date, date_to: date) -> None: with self._connect() as conn: conn.execute( """ INSERT INTO away_periods (telegram_chat_id, date_from, date_to) VALUES (?, ?, ?) """, (telegram_chat_id, date_from.isoformat(), date_to.isoformat()), ) def get_user_by_gitea_login(self, login: str) -> RegisteredUser | None: with self._connect() as conn: row = conn.execute( "SELECT * FROM users WHERE gitea_login = ?", (login,), ).fetchone() if not row: return None return RegisteredUser( telegram_chat_id=int(row["telegram_chat_id"]), telegram_username=row["telegram_username"], email=row["email"], gitea_login=row["gitea_login"], ) def get_user_by_chat_id(self, telegram_chat_id: int) -> RegisteredUser | None: with self._connect() as conn: row = conn.execute( "SELECT * FROM users WHERE telegram_chat_id = ?", (telegram_chat_id,), ).fetchone() if not row: return None return RegisteredUser( telegram_chat_id=int(row["telegram_chat_id"]), telegram_username=row["telegram_username"], email=row["email"], gitea_login=row["gitea_login"], ) def get_available_candidates(self, author_login: str, today: date, allow_self_assign: bool = False) -> list[RegisteredUser]: with self._connect() as conn: if allow_self_assign: rows = conn.execute( """ SELECT u.* FROM users u WHERE NOT EXISTS ( SELECT 1 FROM away_periods a WHERE a.telegram_chat_id = u.telegram_chat_id AND date(a.date_from) <= date(?) AND date(a.date_to) >= date(?) ) """, (today.isoformat(), today.isoformat()), ).fetchall() else: rows = conn.execute( """ SELECT u.* FROM users u WHERE u.gitea_login != ? AND NOT EXISTS ( SELECT 1 FROM away_periods a WHERE a.telegram_chat_id = u.telegram_chat_id AND date(a.date_from) <= date(?) AND date(a.date_to) >= date(?) ) """, (author_login, today.isoformat(), today.isoformat()), ).fetchall() return [ RegisteredUser( telegram_chat_id=int(r["telegram_chat_id"]), telegram_username=r["telegram_username"], email=r["email"], gitea_login=r["gitea_login"], ) for r in rows ] def get_open_reviews_count(self, reviewer_login: str) -> int: with self._connect() as conn: row = conn.execute( "SELECT COUNT(*) as c FROM assignments WHERE reviewer_login = ? AND status = 'open'", (reviewer_login,), ).fetchone() return int(row["c"]) def get_total_reviews_count(self, reviewer_login: str) -> int: with self._connect() as conn: row = conn.execute( "SELECT COUNT(*) as c FROM assignments WHERE reviewer_login = ?", (reviewer_login,), ).fetchone() return int(row["c"]) def add_assignment( self, repo_full_name: str, pr_number: int, author_login: str, reviewer_login: str, ) -> None: with self._connect() as conn: conn.execute( """ INSERT INTO assignments (repo_full_name, pr_number, author_login, reviewer_login, status, created_at) VALUES (?, ?, ?, ?, 'open', ?) ON CONFLICT(repo_full_name, pr_number, reviewer_login) DO UPDATE SET status='open', closed_at=NULL """, (repo_full_name, pr_number, author_login, reviewer_login, utc_now().isoformat()), ) def close_assignments_for_pr(self, repo_full_name: str, pr_number: int) -> None: with self._connect() as conn: conn.execute( """ UPDATE assignments SET status='closed', closed_at=? WHERE repo_full_name=? AND pr_number=? AND status='open' """, (utc_now().isoformat(), repo_full_name, pr_number), ) def get_open_assignments_for_reviewer(self, reviewer_login: str) -> list[sqlite3.Row]: with self._connect() as conn: return conn.execute( """ SELECT repo_full_name, pr_number, author_login, created_at FROM assignments WHERE reviewer_login = ? AND status='open' ORDER BY created_at ASC """, (reviewer_login,), ).fetchall() def get_all_users(self) -> list[RegisteredUser]: with self._connect() as conn: rows = conn.execute("SELECT * FROM users").fetchall() return [ RegisteredUser( telegram_chat_id=int(r["telegram_chat_id"]), telegram_username=r["telegram_username"], email=r["email"], gitea_login=r["gitea_login"], ) for r in rows ] def get_known_repositories(self) -> list[str]: with self._connect() as conn: rows = conn.execute( """ SELECT DISTINCT repo_full_name FROM assignments WHERE repo_full_name IS NOT NULL AND repo_full_name != '' """ ).fetchall() return [str(r["repo_full_name"]) for r in rows] def get_kv(self, key: str) -> str | None: with self._connect() as conn: row = conn.execute("SELECT value FROM kv WHERE key = ?", (key,)).fetchone() return row["value"] if row else None def set_kv(self, key: str, value: str) -> None: with self._connect() as conn: conn.execute( """ INSERT INTO kv (key, value) VALUES (?, ?) ON CONFLICT(key) DO UPDATE SET value=excluded.value """, (key, value), ) def delete_kv(self, key: str) -> None: with self._connect() as conn: conn.execute("DELETE FROM kv WHERE key = ?", (key,)) def get_reminder_time(self, telegram_chat_id: int) -> str | None: return self.get_kv(f"reminder_time:{telegram_chat_id}") def set_reminder_time(self, telegram_chat_id: int, time_str: str) -> None: self.set_kv(f"reminder_time:{telegram_chat_id}", time_str) def get_last_reminder_date(self, telegram_chat_id: int) -> str | None: return self.get_kv(f"last_reminder:{telegram_chat_id}") def set_last_reminder_date(self, telegram_chat_id: int, date_str: str) -> None: self.set_kv(f"last_reminder:{telegram_chat_id}", date_str) def get_assignments_for_pr(self, repo_full_name: str, pr_number: int) -> list[sqlite3.Row]: with self._connect() as conn: return conn.execute( """ SELECT author_login, reviewer_login FROM assignments WHERE repo_full_name = ? AND pr_number = ? AND status = 'open' """, (repo_full_name, pr_number), ).fetchall() class TelegramClient: def __init__(self, token: str) -> None: self.token = token self.base_url = f"https://api.telegram.org/bot{token}" if token else "" async def send_message(self, chat_id: int, text: str) -> None: if not self.base_url: log_warn("[notify] TELEGRAM_BOT_TOKEN пустой, отправка пропущена.") return # По настройке можно отключить отправку уведомлений в выходные (МСК). if not NOTIFY_ON_WEEKENDS and msk_now().weekday() >= 5: log_info(f"[notify] Выходной день, отправка пропущена chat_id={chat_id}") return try: async with httpx.AsyncClient(timeout=20) as client: response = await client.post( f"{self.base_url}/sendMessage", json={"chat_id": chat_id, "text": text}, ) response.raise_for_status() log_info(f"[notify] Сообщение отправлено chat_id={chat_id}") except Exception as exc: # Не роняем обработку webhook, если Telegram API временно недоступен # или указан неверный токен/чат. log_error("telegram send_message", exc) class GiteaClient: def __init__(self, base_url: str, token: str) -> None: self.base_url = base_url self.token = token @property def enabled(self) -> bool: return bool(self.base_url and self.token) async def assign_reviewers(self, repo_full_name: str, pr_number: int, reviewer_logins: list[str]) -> None: if not self.enabled: raise RuntimeError("Gitea клиент не настроен") if not reviewer_logins: return async with httpx.AsyncClient(timeout=20) as client: response = await client.post( f"{self.base_url}/api/v1/repos/{repo_full_name}/pulls/{pr_number}/requested_reviewers", headers={"Authorization": f"token {self.token}"}, json={"reviewers": reviewer_logins}, ) response.raise_for_status() async def add_label(self, repo_full_name: str, pr_number: int, label_name: str) -> None: if not self.enabled: return async with httpx.AsyncClient(timeout=20) as client: # В Gitea метки для PR ставятся через API issue, так как PR является issue. response = await client.post( f"{self.base_url}/api/v1/repos/{repo_full_name}/issues/{pr_number}/labels", headers={"Authorization": f"token {self.token}"}, json={"labels": [label_name]}, ) # Если метки нет, не роняем процесс: для MVP это допустимо. if response.status_code >= 400: return async def get_pull_request(self, repo_full_name: str, pr_number: int) -> dict[str, Any] | None: if not self.enabled: return None async with httpx.AsyncClient(timeout=20) as client: response = await client.get( f"{self.base_url}/api/v1/repos/{repo_full_name}/pulls/{pr_number}", headers={"Authorization": f"token {self.token}"}, ) if response.status_code >= 400: return None return response.json() async def list_open_pull_requests(self, repo_full_name: str, page: int = 1, limit: int = 50) -> list[dict[str, Any]]: if not self.enabled: return [] async with httpx.AsyncClient(timeout=20) as client: response = await client.get( f"{self.base_url}/api/v1/repos/{repo_full_name}/pulls", headers={"Authorization": f"token {self.token}"}, params={"state": "open", "page": page, "limit": limit}, ) if response.status_code >= 400: return [] data = response.json() if not isinstance(data, list): return [] return data async def list_pull_request_reviews(self, repo_full_name: str, pr_number: int) -> list[dict[str, Any]]: if not self.enabled: return [] async with httpx.AsyncClient(timeout=20) as client: response = await client.get( f"{self.base_url}/api/v1/repos/{repo_full_name}/pulls/{pr_number}/reviews", headers={"Authorization": f"token {self.token}"}, ) if response.status_code >= 400: return [] data = response.json() if not isinstance(data, list): return [] return data storage = Storage(DB_PATH) telegram_client = TelegramClient(TELEGRAM_BOT_TOKEN) gitea_client = GiteaClient(GITEA_BASE_URL, GITEA_TOKEN) reminder_task: asyncio.Task | None = None def verify_gitea_signature(raw_body: bytes, signature_header: str | None) -> bool: if not GITEA_WEBHOOK_SECRET: return False if not signature_header: return False signature_value = signature_header.strip() if signature_value.startswith("sha256="): signature_value = signature_value.split("=", 1)[1] digest = hmac.new( GITEA_WEBHOOK_SECRET.encode("utf-8"), raw_body, hashlib.sha256, ).hexdigest() return hmac.compare_digest(digest, signature_value) def choose_reviewer(candidates: list[RegisteredUser]) -> RegisteredUser: return random.choice(candidates) async def build_reviewer_load_counts(candidates: list[RegisteredUser], current_repo_full_name: str) -> dict[str, int]: _ = current_repo_full_name return { candidate.gitea_login: storage.get_total_reviews_count(candidate.gitea_login) for candidate in candidates } async def choose_reviewers_with_remote_load( candidates: list[RegisteredUser], current_repo_full_name: str, reviewers_count: int, ) -> tuple[list[RegisteredUser], dict[str, int]]: counts = await build_reviewer_load_counts(candidates, current_repo_full_name) remaining = list(candidates) selected: list[RegisteredUser] = [] target_count = max(1, min(reviewers_count, len(candidates))) while remaining and len(selected) < target_count: min_count = min(counts[candidate.gitea_login] for candidate in remaining) shortlist = [candidate for candidate in remaining if counts[candidate.gitea_login] == min_count] chosen = choose_reviewer(shortlist) selected.append(chosen) remaining = [candidate for candidate in remaining if candidate.gitea_login != chosen.gitea_login] counts[chosen.gitea_login] += 1 return selected, counts def extract_logins_from_pull_request(pull_request: dict[str, Any] | None) -> set[str]: logins: set[str] = set() if not pull_request: return logins author_login = ((pull_request.get("user") or {}).get("login") or "").strip() if author_login: logins.add(author_login) for key in ("requested_reviewers", "reviewers"): for reviewer in pull_request.get(key) or []: login = ((reviewer or {}).get("login") or "").strip() if login: logins.add(login) return logins def extract_reviewer_logins_from_pull_request(pull_request: dict[str, Any] | None) -> set[str]: logins: set[str] = set() if not pull_request: return logins for key in ("requested_reviewers", "reviewers"): for reviewer in pull_request.get(key) or []: login = ((reviewer or {}).get("login") or "").strip() if login: logins.add(login) return logins def calculate_effective_approvals(reviews: list[dict[str, Any]]) -> int: """Считает актуальное число approve с учетом последнего состояния каждого ревьюера.""" latest_state_by_login: dict[str, tuple[int, str]] = {} for review in reviews: user = review.get("user") or {} login = (user.get("login") or "").strip() if not login: continue state = str(review.get("state") or "").lower() review_id = int(review.get("id") or 0) prev = latest_state_by_login.get(login) if prev is None or review_id >= prev[0]: latest_state_by_login[login] = (review_id, state) return sum(1 for _, state in latest_state_by_login.values() if state == "approved") def get_pr_author_login( repo_full_name: str, pr_number: int, pull_request: dict[str, Any] | None = None ) -> str | None: """Логин создателя (автора) PR.""" if pull_request: login = ((pull_request.get("user") or {}).get("login") or "").strip() if login: return login rows = storage.get_assignments_for_pr(repo_full_name, pr_number) return rows[0]["author_login"] if rows else None async def get_pr_author_chat_id( repo_full_name: str, pr_number: int, pull_request: dict[str, Any] | None = None ) -> int | None: """Chat_id создателя PR в Telegram (если зарегистрирован).""" author_login = get_pr_author_login(repo_full_name, pr_number, pull_request) if not author_login: return None user = storage.get_user_by_gitea_login(author_login) return user.telegram_chat_id if user else None async def get_pr_participant_chat_ids( repo_full_name: str, pr_number: int, pull_request: dict[str, Any] | None = None, exclude_login: str | None = None, ) -> list[int]: """Возвращает chat_id участников PR (автор и ревьюеры), зарегистрированных в боте.""" logins = extract_logins_from_pull_request(pull_request) if len(logins) <= 1: remote_pull = await gitea_client.get_pull_request(repo_full_name, pr_number) logins.update(extract_logins_from_pull_request(remote_pull)) rows = storage.get_assignments_for_pr(repo_full_name, pr_number) for r in rows: logins.add(r["author_login"]) logins.add(r["reviewer_login"]) if exclude_login: logins.discard(exclude_login) chat_ids: list[int] = [] for login in logins: user = storage.get_user_by_gitea_login(login) if user: chat_ids.append(user.telegram_chat_id) else: log_warn(f"[notify] Пользователь '{login}' не зарегистрирован в боте, уведомление пропущено.") return chat_ids async def notify_user_by_login(login: str, text: str) -> None: user = storage.get_user_by_gitea_login(login) if not user: log_warn(f"[notify] Пользователь '{login}' не зарегистрирован в боте, уведомление пропущено.") return await telegram_client.send_message(user.telegram_chat_id, text) async def handle_telegram_command(update: dict[str, Any]) -> None: message = update.get("message") or {} chat = message.get("chat") or {} user = message.get("from") or {} text = (message.get("text") or "").strip() chat_id = chat.get("id") chat_type = chat.get("type") username = user.get("username") if not chat_id or not text.startswith("/"): return if chat_type != "private": log_warn(f"[telegram] Игнорирую команду из не-личного чата id={chat_id}, type={chat_type}") return if not is_user_allowed(int(chat_id), username): log_warn(f"[telegram] Пользователь chat_id={chat_id}, username={username} не входит в allowlist.") return parts = text.split() cmd = parts[0].lower() if cmd in ("/start", "/help"): await telegram_client.send_message( chat_id, ( "Команды:\n" "/register \n" "/away \n" "/reminder [ЧЧ:ММ] — время ежедневного напоминания о ревью\n" "/myreviews" ), ) return if cmd == "/register": if len(parts) != 3: await telegram_client.send_message( chat_id, "Использование: /register ", ) return email = parts[1].lower() gitea_login = parts[2] try: storage.upsert_user( RegisteredUser( telegram_chat_id=int(chat_id), telegram_username=username, email=email, gitea_login=gitea_login, ) ) except sqlite3.IntegrityError: log_warn( f"[telegram] Конфликт регистрации chat_id={chat_id}, " f"email={email}, gitea_login={gitea_login}" ) await telegram_client.send_message( chat_id, "Этот email или gitea_login уже привязан к другому Telegram-аккаунту.", ) return await telegram_client.send_message(chat_id, "Пользователь зарегистрирован.") return if cmd == "/away": if len(parts) != 3: await telegram_client.send_message( chat_id, "Использование: /away ", ) return try: date_from = parse_date(parts[1]) date_to = parse_date(parts[2]) except ValueError: await telegram_client.send_message(chat_id, "Некорректный формат даты.") return if date_to < date_from: await telegram_client.send_message(chat_id, "date_to должна быть больше или равна date_from.") return storage.set_away(int(chat_id), date_from, date_to) await telegram_client.send_message(chat_id, "Период недоступности сохранен.") return if cmd == "/reminder": if len(parts) == 1: t = storage.get_reminder_time(int(chat_id)) msg = f"Время напоминания: {t or REMINDER_TIME} (по умолчанию)" if t else f"Время напоминания: {REMINDER_TIME} (глобальное по умолчанию)" await telegram_client.send_message(chat_id, msg) return if len(parts) != 2: await telegram_client.send_message( chat_id, "Использование: /reminder [ЧЧ:ММ], например /reminder 09:00", ) return time_str = parts[1].strip() if len(time_str) != 5 or time_str[2] != ":": await telegram_client.send_message(chat_id, "Формат времени: ЧЧ:ММ (например 09:00).") return try: h, m = int(time_str[:2]), int(time_str[3:5]) if not (0 <= h <= 23 and 0 <= m <= 59): raise ValueError("out of range") except ValueError: await telegram_client.send_message(chat_id, "Некорректное время. Используйте ЧЧ:ММ (например 09:00).") return storage.set_reminder_time(int(chat_id), time_str) await telegram_client.send_message(chat_id, f"Напоминания будут приходить ежедневно в {time_str}.") return if cmd == "/myreviews": me = storage.get_user_by_chat_id(int(chat_id)) if me is None: log_warn(f"[telegram] /myreviews для незарегистрированного chat_id={chat_id}") await telegram_client.send_message( chat_id, "Учетная запись не найдена. Сначала выполните /register.", ) return rows = storage.get_open_assignments_for_reviewer(me.gitea_login) if not rows: await telegram_client.send_message(chat_id, "Открытых ревью нет.") return lines = ["Открытые ревью:"] for r in rows: lines.append(f"- {r['repo_full_name']}#{r['pr_number']} (автор: {r['author_login']})") await telegram_client.send_message(chat_id, "\n".join(lines)) return async def handle_pr_opened(payload: dict[str, Any]) -> None: pull_request = payload.get("pull_request") or {} repository = payload.get("repository") or {} repo_full_name = repository.get("full_name") pr_number = pull_request.get("number") author_login = (pull_request.get("user") or {}).get("login") if not repo_full_name or not pr_number or not author_login: log_warn("[assign] opened: нет repo/pr/author, событие пропущено.") return if author_login.lower() in AUTO_ASSIGN_EXCLUDED_AUTHORS: log_info( f"[assign] Автоназначение отключено для автора {author_login} " f"в {repo_full_name}#{pr_number}." ) return # В зависимости от версии Gitea ручные ревьюеры могут приходить как # в requested_reviewers, так и в reviewers. existing_reviewer_logins = extract_reviewer_logins_from_pull_request(pull_request) remote_pull = await gitea_client.get_pull_request(repo_full_name, int(pr_number)) if gitea_client.enabled else None existing_reviewer_logins.update(extract_reviewer_logins_from_pull_request(remote_pull)) if existing_reviewer_logins: for reviewer_login in existing_reviewer_logins: storage.add_assignment( repo_full_name=repo_full_name, pr_number=int(pr_number), author_login=author_login, reviewer_login=reviewer_login, ) await notify_user_by_login( reviewer_login, f"Назначено новое ревью: {repo_full_name}#{pr_number} (автор: {author_login})", ) log_info( f"[assign] Ручные ревьюеры уже заданы для {repo_full_name}#{pr_number}: " f"{', '.join(sorted(existing_reviewer_logins))}. Автоназначение пропущено." ) return if not gitea_client.enabled: log_warn("[gitea] Клиент не настроен, автоназначение пропущено.") return candidates = storage.get_available_candidates( author_login=author_login, today=msk_now().date(), allow_self_assign=ALLOW_SELF_ASSIGN, ) if not candidates: log_warn(f"[assign] Нет доступных кандидатов для {repo_full_name}#{pr_number}.") return selected_reviewers, counts = await choose_reviewers_with_remote_load( candidates, repo_full_name, AUTO_ASSIGN_REVIEWERS_COUNT, ) selected_logins = [reviewer.gitea_login for reviewer in selected_reviewers] if not selected_reviewers: log_warn(f"[assign] Не удалось выбрать ревьюеров для {repo_full_name}#{pr_number}.") return log_info( f"[assign] Автоназначение для {repo_full_name}#{pr_number}. " f"Кандидаты={len(candidates)}, выбраны={selected_logins}, нагрузки={dict(counts)}" ) await gitea_client.assign_reviewers(repo_full_name, int(pr_number), selected_logins) await gitea_client.add_label(repo_full_name, int(pr_number), AUTO_ASSIGNED_LABEL) log_info( f"[assign] Запрошено назначение в Gitea для {repo_full_name}#{pr_number}: {selected_logins}; " f"label={AUTO_ASSIGNED_LABEL}" ) for reviewer in selected_reviewers: storage.add_assignment( repo_full_name=repo_full_name, pr_number=int(pr_number), author_login=author_login, reviewer_login=reviewer.gitea_login, ) await telegram_client.send_message( reviewer.telegram_chat_id, f"Назначено новое ревью: {repo_full_name}#{pr_number} (автор: {author_login})", ) log_info( f"[notify] opened для {repo_full_name}#{pr_number}, получателей={len(selected_reviewers)}" ) async def handle_pr_closed(payload: dict[str, Any]) -> None: pull_request = payload.get("pull_request") or {} repository = payload.get("repository") or {} repo_full_name = repository.get("full_name") pr_number = pull_request.get("number") action = payload.get("action") if not repo_full_name or not pr_number: log_warn("[notify] closed/merged: нет repo/pr, событие пропущено.") return participant_chat_ids = await get_pr_participant_chat_ids( repo_full_name, int(pr_number), pull_request=pull_request, ) storage.close_assignments_for_pr(repo_full_name, int(pr_number)) storage.delete_kv(f"merge_ready_notified:{repo_full_name}:{int(pr_number)}") is_merged = action == "merged" or bool(pull_request.get("merged")) status_label = "смержен" if is_merged else "закрыт" title = (pull_request.get("title") or "").strip() or f"#{pr_number}" msg = f"PR {repo_full_name}#{pr_number} {status_label}: {title}" log_info( f"[notify] pr_closed для {repo_full_name}#{pr_number}, " f"status={status_label}, получателей={len(participant_chat_ids)}" ) for chat_id in participant_chat_ids: await telegram_client.send_message(chat_id, msg) async def handle_pr_synchronize(payload: dict[str, Any]) -> None: """Новые коммиты в PR — уведомляем назначенных ревьюеров.""" pull_request = payload.get("pull_request") or {} repository = payload.get("repository") or {} repo_full_name = repository.get("full_name") pr_number = pull_request.get("number") if not repo_full_name or not pr_number: return author_login = (pull_request.get("user") or {}).get("login") or "" msg = f"В PR {repo_full_name}#{pr_number} добавлены новые коммиты (автор: {author_login})." recipient_chat_ids = await get_pr_participant_chat_ids( repo_full_name, int(pr_number), pull_request=pull_request, exclude_login=author_login, ) if not recipient_chat_ids: log_warn( f"[notify] synchronize без получателей для {repo_full_name}#{pr_number}. " f"author={author_login}" ) return log_info( f"[notify] synchronize для {repo_full_name}#{pr_number}, " f"получателей={len(recipient_chat_ids)}" ) for chat_id in recipient_chat_ids: await telegram_client.send_message(chat_id, msg) async def handle_issue_comment(payload: dict[str, Any]) -> None: """Новый комментарий в PR — уведомляем автора и ревьюеров (кроме автора комментария).""" issue = payload.get("issue") or {} if not issue.get("pull_request"): return # это issue, не PR repository = payload.get("repository") or {} repo_full_name = repository.get("full_name") pr_number = issue.get("number") comment = payload.get("comment") or {} commenter_login = (comment.get("user") or {}).get("login") or "" body = (comment.get("body") or "").strip()[:200] if not repo_full_name or not pr_number: return msg = f"Новый комментарий в PR {repo_full_name}#{pr_number}: {body}" recipient_chat_ids = await get_pr_participant_chat_ids( repo_full_name, int(pr_number), exclude_login=commenter_login, ) if not recipient_chat_ids: log_warn( f"[notify] issue_comment без получателей для {repo_full_name}#{pr_number}. " f"commenter={commenter_login}" ) return log_info( f"[notify] issue_comment для {repo_full_name}#{pr_number}, " f"commenter={commenter_login}, получателей={len(recipient_chat_ids)}" ) for chat_id in recipient_chat_ids: await telegram_client.send_message(chat_id, msg) async def handle_pull_request_comment(payload: dict[str, Any]) -> None: """Комментарий в PR из pull_request-события — уведомляем участников, кроме автора комментария.""" pull_request = payload.get("pull_request") or {} repository = payload.get("repository") or {} repo_full_name = repository.get("full_name") pr_number = pull_request.get("number") comment = payload.get("comment") or {} commenter_login = (comment.get("user") or {}).get("login") or "" body = (comment.get("body") or "").strip()[:200] if not repo_full_name or not pr_number or not body: return msg = f"Новый комментарий в PR {repo_full_name}#{pr_number}: {body}" recipient_chat_ids = await get_pr_participant_chat_ids( repo_full_name, int(pr_number), pull_request=pull_request, exclude_login=commenter_login, ) if not recipient_chat_ids: log_warn( f"[notify] pull_request_comment без получателей для {repo_full_name}#{pr_number}. " f"commenter={commenter_login}" ) return log_info( f"[notify] pull_request_comment для {repo_full_name}#{pr_number}, " f"commenter={commenter_login}, получателей={len(recipient_chat_ids)}" ) for chat_id in recipient_chat_ids: await telegram_client.send_message(chat_id, msg) async def handle_pr_reviewed(payload: dict[str, Any]) -> None: """Изменение статуса ревью (approved / request changes) — уведомляем только создателя PR.""" pull_request = payload.get("pull_request") or {} repository = payload.get("repository") or {} review = payload.get("review") or {} repo_full_name = repository.get("full_name") pr_number = pull_request.get("number") if not repo_full_name or not pr_number: log_warn("[notify] reviewed: нет repo/pr, событие пропущено.") return state = (review.get("state") or "").lower() reviewer_login = ( (review.get("user") or {}).get("login") or (payload.get("sender") or {}).get("login") or "" ) if state == "approved": status_text = "одобрен(апрувнут)" elif state in ("request_changes", "request changes"): status_text = "запрошены правки" else: status_text = state or "обновлён" if reviewer_login: msg = f"PR {repo_full_name}#{pr_number}: ревью {status_text} ({reviewer_login})." else: msg = f"PR {repo_full_name}#{pr_number}: ревью {status_text}." author_chat_id = await get_pr_author_chat_id(repo_full_name, int(pr_number), pull_request) if author_chat_id is not None: log_info( f"[notify] reviewed для {repo_full_name}#{pr_number}, " f"state={state or 'unknown'}, reviewer={reviewer_login or '-'}, получателей=1" ) await telegram_client.send_message(author_chat_id, msg) else: log_warn( f"[notify] ревью для {repo_full_name}#{pr_number}: автор не зарегистрирован в боте." ) # Уведомление «можно мерджить» — только создателю PR. if not gitea_client.enabled: return reviews = await gitea_client.list_pull_request_reviews(repo_full_name, int(pr_number)) approvals_count = calculate_effective_approvals(reviews) merge_ready_key = f"merge_ready_notified:{repo_full_name}:{int(pr_number)}" already_notified = storage.get_kv(merge_ready_key) == "1" if approvals_count >= APPROVALS_REQUIRED_FOR_MERGE: if not already_notified and author_chat_id is not None: merge_msg = ( f"PR {repo_full_name}#{pr_number} набрал {approvals_count} approve " f"(порог {APPROVALS_REQUIRED_FOR_MERGE}) — можно мерджить." ) await telegram_client.send_message(author_chat_id, merge_msg) storage.set_kv(merge_ready_key, "1") log_info( f"[notify] merge_ready для {repo_full_name}#{pr_number}, " f"approvals={approvals_count}, threshold={APPROVALS_REQUIRED_FOR_MERGE}" ) else: if already_notified: storage.delete_kv(merge_ready_key) async def reminder_loop() -> None: while True: try: now = msk_now() hh_mm = now.strftime("%H:%M") day_key = now.date().isoformat() users = storage.get_all_users() for user in users: reminder_time = storage.get_reminder_time(user.telegram_chat_id) or REMINDER_TIME if hh_mm != reminder_time: continue already_sent = storage.get_last_reminder_date(user.telegram_chat_id) if already_sent == day_key: continue rows = storage.get_open_assignments_for_reviewer(user.gitea_login) if not rows: continue lines = [f"Ежедневное напоминание ({len(rows)} открыто):"] for row in rows: created = row["created_at"][:10] lines.append( f"- {row['repo_full_name']}#{row['pr_number']} (автор: {row['author_login']}, с {created})" ) await telegram_client.send_message(user.telegram_chat_id, "\n".join(lines)) storage.set_last_reminder_date(user.telegram_chat_id, day_key) await asyncio.sleep(30) except Exception: await asyncio.sleep(5) @asynccontextmanager async def lifespan(_: FastAPI): global reminder_task reminder_task = asyncio.create_task(reminder_loop()) try: yield finally: if reminder_task: reminder_task.cancel() try: await reminder_task except asyncio.CancelledError: pass app = FastAPI(title="CodeReview Bot API", lifespan=lifespan) @app.get("/health") async def health() -> dict[str, str]: return {"status": "ok"} @app.post("/telegram/webhook") async def telegram_webhook( update: dict[str, Any], x_telegram_bot_api_secret_token: str | None = Header(default=None), ) -> dict[str, str]: if not TELEGRAM_BOT_TOKEN: raise HTTPException(status_code=500, detail="Не задан TELEGRAM_BOT_TOKEN") if not TELEGRAM_WEBHOOK_SECRET: raise HTTPException(status_code=500, detail="Не задан TELEGRAM_WEBHOOK_SECRET") if not x_telegram_bot_api_secret_token or not hmac.compare_digest( TELEGRAM_WEBHOOK_SECRET, x_telegram_bot_api_secret_token, ): raise HTTPException(status_code=401, detail="Некорректный telegram webhook secret") try: await handle_telegram_command(update) except Exception as exc: # Telegram повторяет доставку при 5xx, поэтому здесь всегда отвечаем 200 # и логируем причину для отладки. log_error("telegram webhook update handling", exc) return {"ok": "true"} @app.post("/gitea/webhook") async def gitea_webhook( request: Request, x_gitea_event: str | None = Header(default=None), x_gitea_signature: str | None = Header(default=None), ) -> dict[str, str]: if not GITEA_WEBHOOK_SECRET: raise HTTPException(status_code=500, detail="Не задан GITEA_WEBHOOK_SECRET") body = await request.body() if not verify_gitea_signature(body, x_gitea_signature): raise HTTPException(status_code=401, detail="Некорректная подпись") try: payload = json.loads(body.decode("utf-8")) except json.JSONDecodeError as exc: log_error("gitea webhook invalid json", exc) raise HTTPException(status_code=400, detail="Некорректный JSON") from exc log_info(f"[gitea] event={x_gitea_event} action={payload.get('action')}") log_info( "[route] " f"event={x_gitea_event} action={payload.get('action')} " f"has_comment={bool(payload.get('comment'))} " f"has_review={bool(payload.get('review'))} " f"is_pull={bool(payload.get('is_pull'))}" ) try: if x_gitea_event == "issue_comment": log_info("[route] -> handle_issue_comment") await handle_issue_comment(payload) return {"ok": "true"} if x_gitea_event == "pull_request_comment": if payload.get("comment"): log_info("[route] -> handle_pull_request_comment (comment)") await handle_pull_request_comment(payload) elif payload.get("review"): log_info("[route] -> handle_pr_reviewed (from pull_request_comment)") await handle_pr_reviewed(payload) return {"ok": "true"} if x_gitea_event in ("pull_request_review", "pull_request_approved", "pull_request_rejected") and payload.get("review"): await handle_pr_reviewed(payload) return {"ok": "true"} if x_gitea_event == "pull_request": action = payload.get("action") if action == "opened": await handle_pr_opened(payload) elif action in ("commented", "comment"): await handle_pull_request_comment(payload) elif action in ("closed", "merged"): await handle_pr_closed(payload) elif action in ("synchronize", "synchronized"): await handle_pr_synchronize(payload) elif action == "reviewed" and payload.get("review"): await handle_pr_reviewed(payload) return {"ok": "true"} except Exception as exc: # Не падаем на обработчиках событий и логируем безопасно. log_error("gitea webhook event handling", exc) return {"ok": "true"} return {"ok": "ignored"}