уведомления за 10 часов и при истечении

This commit is contained in:
Vlad
2024-10-03 06:52:24 +03:00
parent f42ca88e01
commit d03907b43c
2 changed files with 21 additions and 5 deletions
+16 -2
View File
@@ -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)
+5 -3
View File
@@ -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):
# Проверяем, является ли отправитель администратором