576 lines
21 KiB
Python
576 lines
21 KiB
Python
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"}
|
||
|