diff --git a/README.md b/README.md new file mode 100644 index 0000000..f33863f --- /dev/null +++ b/README.md @@ -0,0 +1,40 @@ +# CodeReview Bot (MVP) + +Реализовано в текущем MVP: + +- Регистрация в Telegram: `/register ` +- Период недоступности: `/away ` +- Webhook Gitea по PR: + - при создании PR без ревьюера бот назначает доступного участника с минимальной нагрузкой + - ставит метку `auto-assigned` (best effort) + - отправляет уведомление назначенному ревьюеру в личные сообщения Telegram +- Ежедневное напоминание в общее время `REMINDER_TIME` (по умолчанию `09:00`) +- Техническая проверка сервиса: `GET /health` + +## Запуск + +```bash +python -m venv .venv +. .venv/Scripts/activate +pip install -r requirements.txt +uvicorn main:app --host 0.0.0.0 --port 8080 --reload +``` + +## Переменные окружения + +- `TELEGRAM_BOT_TOKEN` - обязательно +- `GITEA_BASE_URL` - обязательно для назначения ревьюера, например `https://git.example.com` +- `GITEA_TOKEN` - обязательно для назначения/меток +- `GITEA_WEBHOOK_SECRET` - необязательно, но рекомендуется +- `REMINDER_TIME` - необязательно, по умолчанию `09:00` +- `AUTO_ASSIGNED_LABEL` - необязательно, по умолчанию `auto-assigned` +- `ALLOW_SELF_ASSIGN` - необязательно, по умолчанию `false` (для локального теста можно `true`) +- `BOT_DB_PATH` - необязательно, по умолчанию `botreviewer.sqlite3` + +## Webhook-эндпоинты + +- Обновления Telegram: `POST /telegram/webhook` +- Webhook Gitea: `POST /gitea/webhook` + +Для Telegram укажите URL вашего сервиса с путем `/telegram/webhook`. +Для Gitea включите события `pull_request` и задайте общий секрет. diff --git a/main.py b/main.py new file mode 100644 index 0000000..0e585a3 --- /dev/null +++ b/main.py @@ -0,0 +1,575 @@ +import asyncio +import hashlib +import hmac +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), + ) + + +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) + + +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" + "/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 == "/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") + if not repo_full_name or not pr_number: + return + storage.close_assignments_for_pr(repo_full_name, int(pr_number)) + + +async def reminder_loop() -> None: + while True: + try: + now = utc_now() + hh_mm = now.strftime("%H:%M") + if hh_mm == REMINDER_TIME: + day_key = now.date().isoformat() + already_sent = storage.get_kv("last_reminder_date") + if already_sent != day_key: + users = storage.get_all_users() + for user in users: + 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_kv("last_reminder_date", 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 = await request.json() + if x_gitea_event != "pull_request": + return {"ok": "ignored"} + action = payload.get("action") + if action == "opened": + await handle_pr_opened(payload) + elif action in ("closed", "merged"): + await handle_pr_closed(payload) + return {"ok": "true"} + diff --git a/requirements.txt b/requirements.txt new file mode 100644 index 0000000..d23d558 --- /dev/null +++ b/requirements.txt @@ -0,0 +1,3 @@ +fastapi +uvicorn +httpx