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, 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") ALLOW_SELF_ASSIGN = os.getenv("ALLOW_SELF_ASSIGN", "false").lower() == "true" @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 parse_date(value: str) -> date: return datetime.strptime(value, "%Y-%m-%d").date() 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 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_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 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: 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() except Exception as exc: # Не роняем обработку webhook, если Telegram API временно недоступен # или указан неверный токен/чат. print(f"[telegram] Ошибка отправки сообщения: {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_reviewer(self, repo_full_name: str, pr_number: int, reviewer_login: str) -> None: if not self.enabled: raise RuntimeError("Gitea клиент не настроен") 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_login]}, ) 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 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 True if not signature_header: return False digest = hmac.new( GITEA_WEBHOOK_SECRET.encode("utf-8"), raw_body, hashlib.sha256, ).hexdigest() return hmac.compare_digest(digest, signature_header) def choose_reviewer(candidates: list[RegisteredUser]) -> RegisteredUser: counts: dict[str, int] = {} for candidate in candidates: counts[candidate.gitea_login] = storage.get_open_reviews_count(candidate.gitea_login) min_count = min(counts.values()) shortlist = [c for c in candidates if counts[c.gitea_login] == min_count] return random.choice(shortlist) def get_pr_participant_chat_ids( repo_full_name: str, pr_number: int, exclude_login: str | None = None ) -> list[int]: """Возвращает chat_id участников PR (автор и ревьюеры), зарегистрированных в боте. exclude_login не включается.""" rows = storage.get_assignments_for_pr(repo_full_name, pr_number) if not rows: return [] logins: set[str] = set() 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) return chat_ids 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") username = user.get("username") if not chat_id or not text.startswith("/"): 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] storage.upsert_user( RegisteredUser( telegram_chat_id=int(chat_id), telegram_username=username, email=email, gitea_login=gitea_login, ) ) 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: 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 {} requested_reviewers = pull_request.get("requested_reviewers") or [] if requested_reviewers: return 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: return if not gitea_client.enabled: return candidates = storage.get_available_candidates( author_login=author_login, today=utc_now().date(), allow_self_assign=ALLOW_SELF_ASSIGN, ) if not candidates: return reviewer = choose_reviewer(candidates) await gitea_client.assign_reviewer(repo_full_name, int(pr_number), reviewer.gitea_login) await gitea_client.add_label(repo_full_name, int(pr_number), AUTO_ASSIGNED_LABEL) 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})", ) 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: return participant_chat_ids = get_pr_participant_chat_ids(repo_full_name, int(pr_number)) storage.close_assignments_for_pr(repo_full_name, int(pr_number)) status_label = "смержен" if action == "merged" else "закрыт" title = (pull_request.get("title") or "").strip() or f"#{pr_number}" msg = f"PR {repo_full_name}#{pr_number} {status_label}: {title}" 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 rows = storage.get_assignments_for_pr(repo_full_name, int(pr_number)) author_login = (pull_request.get("user") or {}).get("login") or "" msg = f"В PR {repo_full_name}#{pr_number} добавлены новые коммиты (автор: {author_login})." for r in rows: reviewer_login = r["reviewer_login"] user = storage.get_user_by_gitea_login(reviewer_login) if user: await telegram_client.send_message(user.telegram_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}" for chat_id in get_pr_participant_chat_ids(repo_full_name, int(pr_number), exclude_login=commenter_login): await telegram_client.send_message(chat_id, msg) async def handle_pr_reviewed(payload: dict[str, Any]) -> None: """Изменение статуса ревью (approved / request changes) — уведомляем участников.""" 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: return state = (review.get("state") or "").lower() reviewer_login = (review.get("user") or {}).get("login") or "" if state == "approved": status_text = "одобрен" elif state in ("request_changes", "request changes"): status_text = "запрошены правки" else: status_text = state or "обновлён" msg = f"PR {repo_full_name}#{pr_number}: ревью {status_text} ({reviewer_login})." for chat_id in get_pr_participant_chat_ids(repo_full_name, int(pr_number), exclude_login=reviewer_login): await telegram_client.send_message(chat_id, msg) async def reminder_loop() -> None: while True: try: now = utc_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]) -> dict[str, str]: if not TELEGRAM_BOT_TOKEN: raise HTTPException(status_code=500, detail="Не задан TELEGRAM_BOT_TOKEN") try: await handle_telegram_command(update) except Exception as exc: # Telegram повторяет доставку при 5xx, поэтому здесь всегда отвечаем 200 # и логируем причину для отладки. print(f"[telegram] Ошибка обработки апдейта: {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]: body = await request.body() if not verify_gitea_signature(body, x_gitea_signature): raise HTTPException(status_code=401, detail="Некорректная подпись") payload = json.loads(body.decode("utf-8")) if x_gitea_event == "issue_comment": await handle_issue_comment(payload) return {"ok": "true"} if x_gitea_event == "pull_request_review" 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 ("closed", "merged"): await handle_pr_closed(payload) elif action == "synchronize": await handle_pr_synchronize(payload) elif action == "reviewed" and payload.get("review"): await handle_pr_reviewed(payload) return {"ok": "true"} return {"ok": "ignored"}