1
0
Fork 0
mirror of https://github.com/EDeev/chatping_abobot.git synced 2026-10-07 20:49:45 +03:00
chatping_abobot/code/db.py
Egor Deev a2a0adcf8e AboBot 4.0: PostgreSQL, обработчики по роутерам, HTML-разметка, новые команды
База — 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.
2026-10-06 11:16:00 +00:00

199 lines
9 KiB
Python
Raw Permalink 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.

"""Работа с 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]