ReviewBot/main.py
2026-02-20 21:03:38 +03:00

576 lines
21 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

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 <git_email> <gitea_login>\n"
"/away <YYYY-MM-DD> <YYYY-MM-DD>\n"
"/myreviews"
),
)
return
if cmd == "/register":
if len(parts) != 3:
await telegram_client.send_message(
chat_id,
"Использование: /register <git_email> <gitea_login>",
)
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 <YYYY-MM-DD> <YYYY-MM-DD>",
)
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"}