mirror of
https://github.com/EDeev/chatping_abobot.git
synced 2026-10-07 20:49:45 +03:00
База — PostgreSQL (asyncpg) вместо четырёх файлов SQLite с таблицей на каждый чат: chats, users, members, chat_stats и member_stats с полем «период» — месяцы хранятся, а не обнуляются; счётчики растут одним запросом без блокировки бота. scripts/migrate_sqlite.py переносит старые базы со сверкой сумм (проверено на копии рабочих баз: 20 694 пользователя, 45 чатов, 3 142 038 сообщений — совпадает). Учёт бота в чатах: статус и права бота (событие my_chat_member и сверка с Telegram при запуске и раз в 6 часов), история изменений в bot_status_history; из чатов, где бота исключили, данные не удаляются — меняется только статус. Большие чаты (больше 100 участников): /all — только для админов, раз в 5 минут и только писавшие за 30 дней; упоминание одного человека по имени — не чаще раза в минуту. В любом чате /all делится на сообщения до 4096 символов (раньше в чате на 14 тысяч падал). Новое: /settings (упоминания, имена в голосовых, ивенты, удаление служебных сообщений — меняют админы), /top, /month (итоги месяца по команде). /stop_bot раньше ничего не выключал — флаг не проверялся; теперь работает. Участник, вышедший из чата, больше не упоминается, но его статистика сохраняется. Ошибки — в технический чат DEBUG_CHAT_ID. Код: handlers.py разбит на роутеры, Markdown заменён на HTML (экранирование через html.escape), тесты на PostgreSQL (18), CI на Python 3.10 и 3.12.
199 lines
9 KiB
Python
199 lines
9 KiB
Python
"""Работа с PostgreSQL. Счётчики увеличиваются одним запросом (INSERT … ON CONFLICT DO UPDATE)"""
|
||
import json
|
||
import os
|
||
from datetime import date
|
||
|
||
import asyncpg
|
||
|
||
COLUMNS = {1: "mes", 2: "rep", 3: "com", 4: "url", 5: "med", 6: "sti", 7: "voi"}
|
||
STAT_FIELDS = ("mes", "rep", "com", "url", "med", "sti", "voi")
|
||
SETTINGS = ("events_enabled", "mentions_enabled", "voice_names", "delete_service")
|
||
|
||
pool: asyncpg.Pool | None = None
|
||
|
||
|
||
def month_key(day=None):
|
||
return (day or date.today()).strftime("%Y-%m")
|
||
|
||
|
||
async def connect(dsn):
|
||
global pool
|
||
pool = await asyncpg.create_pool(dsn, min_size=1, max_size=5)
|
||
with open(os.path.join(os.path.dirname(__file__), "schema.sql"), encoding="utf-8") as f:
|
||
await pool.execute(f.read())
|
||
return pool
|
||
|
||
|
||
async def close():
|
||
if pool:
|
||
await pool.close()
|
||
|
||
|
||
# ЧАТЫ И УЧАСТНИКИ
|
||
async def ensure_chat(chat_id, title=None):
|
||
await pool.execute(
|
||
"INSERT INTO chats (id, title) VALUES ($1, $2) "
|
||
"ON CONFLICT (id) DO UPDATE SET title = COALESCE(EXCLUDED.title, chats.title)", chat_id, title)
|
||
|
||
|
||
async def chat_exists(chat_id):
|
||
return await pool.fetchval("SELECT EXISTS (SELECT 1 FROM chats WHERE id = $1)", chat_id)
|
||
|
||
|
||
async def migrate_chat(old_id, new_id):
|
||
"""Группа стала супергруппой — у неё новый id; связанные таблицы обновятся каскадом"""
|
||
await pool.execute("UPDATE chats SET id = $2 WHERE id = $1 AND NOT EXISTS (SELECT 1 FROM chats WHERE id = $2)",
|
||
old_id, new_id)
|
||
|
||
|
||
async def ensure_user(user_id):
|
||
await pool.execute("INSERT INTO users (id) VALUES ($1) ON CONFLICT DO NOTHING", user_id)
|
||
|
||
|
||
async def get_setting(chat_id, name):
|
||
assert name in SETTINGS
|
||
value = await pool.fetchval(f"SELECT {name} FROM chats WHERE id = $1", chat_id)
|
||
return True if value is None else value
|
||
|
||
|
||
async def toggle_setting(chat_id, name):
|
||
assert name in SETTINGS
|
||
return await pool.fetchval(f"UPDATE chats SET {name} = NOT {name} WHERE id = $1 RETURNING {name}", chat_id)
|
||
|
||
|
||
async def settings(chat_id):
|
||
row = await pool.fetchrow(f"SELECT {', '.join(SETTINGS)} FROM chats WHERE id = $1", chat_id)
|
||
return dict(row) if row else dict.fromkeys(SETTINGS, True)
|
||
|
||
|
||
async def chat_names(chat_id):
|
||
"""Имена действующих участников: {имя: id}"""
|
||
rows = await pool.fetch("SELECT name, user_id FROM members WHERE chat_id = $1 AND left_at IS NULL AND name <> ''",
|
||
chat_id)
|
||
return {r["name"]: r["user_id"] for r in rows}
|
||
|
||
|
||
async def active_members(chat_id, days=None):
|
||
"""Действующие участники: [(id, имя)]; days — только писавшие за последние N дней"""
|
||
query = "SELECT user_id, name FROM members WHERE chat_id = $1 AND left_at IS NULL AND name <> ''"
|
||
args = [chat_id]
|
||
if days:
|
||
query += " AND last_message_at > now() - make_interval(days => $2)"
|
||
args.append(days)
|
||
return [(r["user_id"], r["name"]) for r in await pool.fetch(query + " ORDER BY name", *args)]
|
||
|
||
|
||
async def member_count(chat_id):
|
||
return await pool.fetchval("SELECT count(*) FROM members WHERE chat_id = $1 AND left_at IS NULL", chat_id)
|
||
|
||
|
||
async def member_name(chat_id, user_id):
|
||
return await pool.fetchval("SELECT name FROM members WHERE chat_id = $1 AND user_id = $2", chat_id, user_id)
|
||
|
||
|
||
async def set_custom_name(chat_id, user_id, name):
|
||
await pool.execute("UPDATE members SET name = $3, custom_name = TRUE WHERE chat_id = $1 AND user_id = $2",
|
||
chat_id, user_id, name)
|
||
|
||
|
||
async def reset_custom_name(chat_id, user_id, name):
|
||
"""Возвращает True, если своё имя было задано"""
|
||
return await pool.fetchval(
|
||
"UPDATE members SET name = $3, custom_name = FALSE WHERE chat_id = $1 AND user_id = $2 AND custom_name "
|
||
"RETURNING TRUE", chat_id, user_id, name) or False
|
||
|
||
|
||
async def member_left(chat_id, user_id):
|
||
await pool.execute("UPDATE members SET left_at = now() WHERE chat_id = $1 AND user_id = $2", chat_id, user_id)
|
||
|
||
|
||
async def mark_all(chat_id):
|
||
await pool.execute("UPDATE chats SET last_all_at = now() WHERE id = $1", chat_id)
|
||
|
||
|
||
async def last_all(chat_id):
|
||
return await pool.fetchval("SELECT last_all_at FROM chats WHERE id = $1", chat_id)
|
||
|
||
|
||
# СТАТИСТИКА
|
||
async def count(chat_id, user_id, var_ids, name, title=None, day=None):
|
||
"""Учёт события: участник появляется в чате (имя обновляется, если не задано своё),
|
||
счётчики var_ids растут за всё время и за текущий месяц — одной транзакцией"""
|
||
periods = ("all", month_key(day))
|
||
columns = [COLUMNS[v] for v in var_ids]
|
||
increments = ", ".join(f"{c} = s.{c} + EXCLUDED.{c}" for c in columns)
|
||
values = ", ".join("1" for _ in columns)
|
||
column_list = ", ".join(columns)
|
||
|
||
async with pool.acquire() as conn, conn.transaction():
|
||
await conn.execute(
|
||
"INSERT INTO chats (id, title) VALUES ($1, $2) "
|
||
"ON CONFLICT (id) DO UPDATE SET title = COALESCE(EXCLUDED.title, chats.title)", chat_id, title)
|
||
await conn.execute("INSERT INTO users (id) VALUES ($1) ON CONFLICT DO NOTHING", user_id)
|
||
await conn.execute(
|
||
"INSERT INTO members (chat_id, user_id, name, last_message_at) VALUES ($1, $2, $3, now()) "
|
||
"ON CONFLICT (chat_id, user_id) DO UPDATE SET last_message_at = now(), left_at = NULL, "
|
||
"name = CASE WHEN members.custom_name THEN members.name ELSE EXCLUDED.name END",
|
||
chat_id, user_id, name)
|
||
for period in periods:
|
||
await conn.execute(
|
||
f"INSERT INTO chat_stats AS s (chat_id, period, {column_list}) VALUES ($1, $2, {values}) "
|
||
f"ON CONFLICT (chat_id, period) DO UPDATE SET {increments}", chat_id, period)
|
||
await conn.execute(
|
||
f"INSERT INTO member_stats AS s (chat_id, user_id, period, {column_list}) VALUES ($1, $2, $3, {values}) "
|
||
f"ON CONFLICT (chat_id, user_id, period) DO UPDATE SET {increments}", chat_id, user_id, period)
|
||
|
||
|
||
def _stats(row):
|
||
return tuple(row[f] for f in STAT_FIELDS) if row else (0,) * len(STAT_FIELDS)
|
||
|
||
|
||
async def chat_stats(chat_id, period):
|
||
return _stats(await pool.fetchrow("SELECT * FROM chat_stats WHERE chat_id = $1 AND period = $2", chat_id, period))
|
||
|
||
|
||
async def member_stats(chat_id, user_id, period):
|
||
return _stats(await pool.fetchrow("SELECT * FROM member_stats WHERE chat_id = $1 AND user_id = $2 AND period = $3",
|
||
chat_id, user_id, period))
|
||
|
||
|
||
async def top(chat_id, period, field="mes", limit=10):
|
||
"""Лучшие участники по счётчику: [(id, имя, значение)]"""
|
||
assert field in STAT_FIELDS
|
||
rows = await pool.fetch(
|
||
f"SELECT s.user_id, m.name, s.{field} AS value FROM member_stats s "
|
||
f"JOIN members m USING (chat_id, user_id) WHERE s.chat_id = $1 AND s.period = $2 AND s.{field} > 0 "
|
||
f"ORDER BY s.{field} DESC, m.name LIMIT $3", chat_id, period, limit)
|
||
return [(r["user_id"], r["name"], r["value"]) for r in rows]
|
||
|
||
|
||
# БОТ В ЧАТАХ
|
||
IN_CHAT = ("member", "administrator", "restricted", "creator")
|
||
|
||
|
||
async def set_bot_status(chat_id, status, rights, source, member_count=None, chat_type=None, title=None):
|
||
"""Записывает статус бота в чате; в историю — только если статус или права изменились"""
|
||
rights_json = None if rights is None else json.dumps(rights, sort_keys=True)
|
||
async with pool.acquire() as conn, conn.transaction():
|
||
await conn.execute(
|
||
"INSERT INTO chats (id, title) VALUES ($1, $2) "
|
||
"ON CONFLICT (id) DO UPDATE SET title = COALESCE(EXCLUDED.title, chats.title)", chat_id, title)
|
||
changed = await conn.fetchval(
|
||
"SELECT bot_status IS DISTINCT FROM $2 OR bot_rights IS DISTINCT FROM $3::jsonb FROM chats WHERE id = $1",
|
||
chat_id, status, rights_json)
|
||
await conn.execute(
|
||
"UPDATE chats SET bot_status = $2, bot_rights = $3::jsonb, checked_at = now(), "
|
||
"bot_status_at = CASE WHEN $4 THEN now() ELSE bot_status_at END, "
|
||
"member_count = COALESCE($5, member_count), chat_type = COALESCE($6, chat_type) WHERE id = $1",
|
||
chat_id, status, rights_json, changed, member_count, chat_type)
|
||
if changed:
|
||
await conn.execute("INSERT INTO bot_status_history (chat_id, status, rights, source) "
|
||
"VALUES ($1, $2, $3::jsonb, $4)", chat_id, status, rights_json, source)
|
||
return changed
|
||
|
||
|
||
async def chats_to_check():
|
||
"""Чаты для сверки: все, где бот считается состоящим, и те, где статус неизвестен"""
|
||
rows = await pool.fetch("SELECT id FROM chats WHERE bot_status IS NULL OR bot_status = ANY($1::text[])",
|
||
list(IN_CHAT))
|
||
return [r["id"] for r in rows]
|