diff --git a/database.py b/database.py index dca8bd56..dc5fe3cf 100644 --- a/database.py +++ b/database.py @@ -6,6 +6,7 @@ from config import DATABASE_URL async def init_db(): conn = await asyncpg.connect(DATABASE_URL) + # Создаем таблицу connections, если она не существует await conn.execute(''' CREATE TABLE IF NOT EXISTS connections ( @@ -14,6 +15,7 @@ async def init_db(): trial INTEGER NOT NULL DEFAULT 0 ) ''') + # Создаем таблицу keys, если она не существует await conn.execute(''' CREATE TABLE IF NOT EXISTS keys ( @@ -23,10 +25,12 @@ async def init_db(): created_at BIGINT NOT NULL, expiry_time BIGINT NOT NULL, key TEXT NOT NULL, - server_id TEXT NOT NULL DEFAULT 'server1', -- новое поле для идентификатора сервера + server_id TEXT NOT NULL DEFAULT 'server1', -- поле для идентификатора сервера + notified BOOLEAN NOT NULL DEFAULT FALSE, -- новое поле для статуса уведомления PRIMARY KEY (tg_id, client_id) ) ''') + # Добавляем поле server_id в таблицу keys, если его нет try: await conn.execute(''' @@ -36,8 +40,18 @@ async def init_db(): except asyncpg.exceptions.DuplicateColumnError: # Если поле уже существует, ничего не делаем pass - await conn.close() + + # Добавляем поле notified в таблицу keys, если его нет + try: + await conn.execute(''' + ALTER TABLE keys + ADD COLUMN notified BOOLEAN NOT NULL DEFAULT FALSE + ''') + except asyncpg.exceptions.DuplicateColumnError: + # Если поле уже существует, ничего не делаем + pass + await conn.close() async def add_connection(tg_id: int, balance: float = 0.0, trial: int = 0): conn = await asyncpg.connect(DATABASE_URL) diff --git a/handlers/notifications.py b/handlers/notifications.py index 4a7d408c..a4fac2d8 100644 --- a/handlers/notifications.py +++ b/handlers/notifications.py @@ -21,11 +21,11 @@ async def notify_expiring_keys(bot: Bot): try: conn = await asyncpg.connect(DATABASE_URL) try: - # Получаем все ключи, которые истекают в течение следующих 10 часов + # Получаем все ключи, которые истекают в течение следующих 10 часов и еще не были уведомлены threshold_time = (datetime.utcnow() + timedelta(hours=10)).timestamp() * 1000 # В миллисекундах records = await conn.fetch(''' SELECT tg_id, email, expiry_time, client_id, server_id FROM keys - WHERE expiry_time <= $1 AND expiry_time > $2 + WHERE expiry_time <= $1 AND expiry_time > $2 AND notified = FALSE ''', threshold_time, datetime.utcnow().timestamp() * 1000) for record in records: @@ -51,6 +51,9 @@ async def notify_expiring_keys(bot: Bot): try: await bot.send_message(chat_id=tg_id, text=message, parse_mode='HTML', reply_markup=keyboard) + + # Обновляем статус уведомления в базе данных + await conn.execute('UPDATE keys SET notified = TRUE WHERE client_id = $1', record['client_id']) except Exception as e: print(f"Ошибка при отправке сообщения пользователю {tg_id}: {e}. Пропускаем этого пользователя.") @@ -92,7 +95,6 @@ async def notify_expiring_keys(bot: Bot): except Exception as e: print(f"Ошибка при отправке уведомлений: {e}") - @router.message(Command('send_to_all')) async def send_message_to_all_clients(message: types.Message): # Проверяем, является ли отправитель администратором