ReviewBot/main.py
Raykov-MS dbe3a7c34d wda
2026-02-28 21:27:06 +03:00

1006 lines
39 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 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
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()
}
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_info(msg: str) -> None:
print(f"[info] {msg}")
def log_warn(msg: str) -> None:
print(f"[warn] {msg}")
def log_error(context: str, exc: Exception) -> None:
# Не логируем str(exc), чтобы не утекали токены/URL/чувствительные детали.
print(f"[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 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 временно недоступен
# или указан неверный токен/чат.
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
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
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 <git_email> <gitea_login>\n"
"/away <YYYY-MM-DD> <YYYY-MM-DD>\n"
"/reminder [ЧЧ:ММ] — время ежедневного напоминания о ревью\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]
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 <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 == "/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:
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)
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})",
)
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 = 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))
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}"
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})."
for chat_id in await get_pr_participant_chat_ids(
repo_full_name,
int(pr_number),
pull_request=pull_request,
exclude_login=author_login,
):
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}"
for chat_id in await 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_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}"
for chat_id in await get_pr_participant_chat_ids(
repo_full_name,
int(pr_number),
pull_request=pull_request,
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 await get_pr_participant_chat_ids(
repo_full_name,
int(pr_number),
pull_request=pull_request,
exclude_login=reviewer_login,
):
await telegram_client.send_message(chat_id, msg)
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')}")
try:
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 ("commented", "comment"):
await handle_pull_request_comment(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"}
except Exception as exc:
# Не падаем на обработчиках событий и логируем безопасно.
log_error("gitea webhook event handling", exc)
return {"ok": "true"}
return {"ok": "ignored"}