Add files via upload

This commit is contained in:
imletbruh
2026-05-31 18:24:39 +05:00
committed by GitHub
parent fcaa53b5ce
commit 9de1a40f1f
6 changed files with 1392 additions and 1227 deletions
+124 -178
View File
@@ -1,27 +1,41 @@
# 🚀 QIYANA AUTO-BUMP BOT для Lolz.live # 🚀 QIYANA AUTO-BUMP BOT для Lolz.live
Production-ready Telegram бот для автоматического поднятия тем на форуме Lolz.live с современной архитектурой, batch API и полной типизацией. Telegram бот для автоматического поднятия тем на форуме Lolz.live. Batch API, динамические настройки (меняются на лету без перезапуска), персистентная статистика.
## ✨ Возможности ## ✨ Возможности
- **Добавление тем** - через запятую (12345, 67890, 11111) с batch API - **Добавление тем** через запятую, названия загружаются сразу через batch API
- 🗑️ **Удаление тем** - через интерактивное меню - 🗑️ **Удаление тем** через интерактивное меню
- 📋 **Список тем** - с датой последнего поднятия - 📋 **Список тем** с названиями из БД и датой последнего поднятия
- 🚀 **Ручное поднятие** - поднять все темы немедленно через batch API - 🚀 **Ручное поднятие** все темы немедленно через batch API
- **Автоподнятие** - каждые N часов автоматически - 🔄 **Обновление названий** — синхронизация названий тем с форумом по кнопке
- 📊 **Статистика** - успешность, количество поднятий, uptime - **Автоподнятие** — каждые N часов автоматически
- 🔔 **Уведомления** - о каждом поднятии темы - 🛠️ **Динамические настройки** — интервал, batch size, автобамп меняются через меню без перезапуска
- 🛡️ **Обработка ошибок** - rate limits, network failures, invalid tokens - 📊 **Статистика** — сохраняется в БД, не теряется при перезапуске
- **Batch API** - до 10 тем за один запрос (10x быстрее) - 🔔 **Уведомления** — о результате каждого поднятия
- 🛡️ **Защита доступа** — ботом управляет только админ (по Telegram ID)
-**Batch API** — до 10 тем за один запрос
## 🎮 Интерфейс ## 🎮 Интерфейс
``` ```
┌─────────────────────────────────┐ ┌──────────────────────────────────────
Добавить темы │ 📋 Список тем Add topics │ 📋 List of topics
│ 🗑️ Удалить тему │ 🚀 Поднять темы │ 🗑️ Delete topic │ 🚀 Bump topics
📊 Статистика │ 👤 Автор 🔄 Refresh │ 📊 Statistics
└─────────────────────────────────┘ │ 🛠️ Settings │ 👤 Author │
└──────────────────────────────────────┘
```
### Меню настроек
```
┌──────────────────────────────────┐
│ ⏰ Set Interval │
│ 📦 Set Batch Size │
│ 🔄 Toggle Auto-Bump │
│ ↩️ Back to Menu │
└──────────────────────────────────┘
``` ```
## 📦 Установка ## 📦 Установка
@@ -29,9 +43,9 @@ Production-ready Telegram бот для автоматического подн
### Требования ### Требования
- Python 3.12+ - Python 3.12+
- pip или uv - pip
### 1. Установите зависимости ### 1. Клонируйте и установите зависимости
```bash ```bash
pip install -r requirements.txt pip install -r requirements.txt
@@ -49,29 +63,32 @@ pip install -r requirements.txt
2. Создайте токен с правами `read`, `post` 2. Создайте токен с правами `read`, `post`
3. Скопируйте токен 3. Скопируйте токен
### 3. Настройте config.json **Ваш Telegram ID:**
1. Напишите [@userinfobot](https://t.me/userinfobot) и получите свой ID
```json ### 3. Настройте `.env`
{
"bot": { ```env
"api_token": "1234567890:ABCdefGHIjklMNOpqrsTUVwxyz", # Telegram Bot
"img_url": "https://wallpapers-clan.com/wp-content/uploads/2024/04/dark-anime-girl-with-red-eyes-desktop-wallpaper-preview.jpg", BOT_API_TOKEN=1234567890:ABCdefGHIjklMNOpqrsTUVwxyz
"author_url": "https://lolz.live/kqlol/" BOT_IMG_URL=https://wallpapers-clan.com/wp-content/uploads/2024/04/dark-anime-girl-with-red-eyes-desktop-wallpaper-preview.jpg
}, BOT_AUTHOR_URL=https://lolz.live/kqlol/
"api": {
"base_url": "https://prod-api.lolz.live", # Lolz API
"auth_token": "eyJ0eXAiOiJKV1QiLCJhbGc...", API_BASE_URL=https://prod-api.lolz.live
"batch_size": 10 API_AUTH_TOKEN=eyJ0eXAiOiJKV1QiLCJhbGc...
}, API_BATCH_SIZE=10
"database": {
"path": "threads.db" # Database
}, DB_PATH=threads.db
"scheduling": {
"bump_interval_hours": 12, # Scheduling
"bump_delay_seconds": 2, BUMP_INTERVAL_HOURS=12
"enable_auto_bump": true BUMP_DELAY_SECONDS=2
} ENABLE_AUTO_BUMP=true
}
# Admin (ваш Telegram user ID)
ADMIN_USER_ID=123456789
``` ```
### 4. Запустите бота ### 4. Запустите бота
@@ -84,210 +101,139 @@ python app.py
### Добавление тем ### Добавление тем
1. Нажмите ** Добавить темы** 1. Нажмите ** Add topics**
2. Введите ID через запятую: `12345, 67890, 11111` 2. Введите ID через запятую: `12345, 67890, 11111`
3. Бот автоматически получит названия тем через batch API (до 10 за запрос) 3. Бот сразу загрузит реальные названия через batch API и сохранит в БД
**Пример:** ### Обновление названий
```
Ввод: 9247920, 9247921, 9247922
Результат: 1. Нажмите **🔄 Refresh**
✅ 9247920 - Qiyanas steam idler 2. Бот загрузит актуальные названия всех тем через batch API
✅ 9247921 - Discord bot 3. Обновлённые названия сохранятся в БД
✅ 9247922 - VPN service
```
**Batch API преимущества:** Полезно, если темы переименовали на форуме — не нужно удалять и добавлять заново.
- Добавление 10 тем = 1 batch запрос вместо 10 обычных
- Добавление 25 тем = 3 batch запроса вместо 25 обычных
- Экономия API лимитов в 10 раз
### Удаление темы ### Удаление темы
1. Нажмите **🗑️ Удалить тему** 1. Нажмите **🗑️ Delete topic**
2. Выберите тему из списка 2. Выберите тему из списка
3. Подтвердите удаление
### Ручное поднятие ### Ручное поднятие
1. Нажмите **🚀 Поднять темы** 1. Нажмите **🚀 Bump topics**
2. Бот поднимет все темы через batch API (до 10 за запрос) 2. Бот поднимет все темы через batch API
3. Получите уведомление о каждой теме: 3. Получите уведомление о каждой теме:
``` ```
[1/3] ✅ Тема 9247920 поднята успешно [1/3] ✅ Тема 9247920 поднята успешно
[2/3] ✅ Тема 9247921 поднята успешно [2/3] ✅ Тема 9247921 поднята успешно
[3/3] Тема 9247922 поднята успешно [3/3] Тема 9247922: нужно подождать
``` ```
**Batch API преимущества:**
- Поднятие 10 тем = 1 batch запрос (0.1 * 10 = 1 batch)
- Поднятие 50 тем = 5 batch запросов вместо 50 обычных
- Скорость выполнения увеличена в ~10 раз
### Просмотр списка
Нажмите **📋 Список тем** - покажет:
- ID темы
- Название темы
- Дату последнего поднятия
### Статистика ### Статистика
Нажмите **📊 Статистика** - покажет: Нажмите **📊 Statistics** покажет:
- Количество тем - Количество тем (всего / готовы к поднятию)
- Готовые к поднятию - Статистика поднятий (всего / успешно / %)
- Всего попыток - Дата последнего бампа
- Успешность (%) - Текущие настройки
- Настройки интервала
- Время работы бота - Время работы бота
## ⚙️ Настройки Статистика сохраняется в SQLite и не сбрасывается при перезапуске.
### Интервал автоподнятия ## ⚙️ Настройки (меняются на лету)
В `config.json` измените `bump_interval_hours`: Все настройки изменяются через кнопку **🛠️ Settings** без остановки бота:
| Значение | Интервал | | Кнопка | Что меняет | Применяется |
|----------|----------| |--------|-----------|-------------|
| `12` | 12 часов | | ⏰ Set Interval | Интервал автобампа (часы) | Мгновенно, перезапускает цикл |
| `6` | 6 часов | | 📦 Set Batch Size | Размер batch (1–10) | Мгновенно для следующих запросов |
| `24` | 24 часа | | 🔄 Toggle Auto-Bump | Включить/выключить автобамп | Мгновенно |
### Размер batch запросов Настройки хранятся в БД и сохраняются между перезапусками.
В `config.json` измените `batch_size`:
| Значение | Описание |
|----------|----------|
| `10` | Максимум (рекомендуется) - 10 тем за запрос |
| `5` | Средний - 5 тем за запрос |
| `1` | Отключить batch - по 1 теме |
**Примечание:** Batch API Lolz.live поддерживает до 10 запросов в одном batch. Каждый запрос = 0.1 batch, итого 10 запросов = 1 полный batch.
### Тестовый режим (5 минут)
```json
"scheduling": {
"bump_interval_hours": 0.0833,
"bump_delay_seconds": 1,
"enable_auto_bump": true
}
```
## 🔄 Как работает автоподнятие ## 🔄 Как работает автоподнятие
1. Запускаете бота 1. Запускаете бота
2. Вручную поднимаете темы через 🚀 (первый раз) 2. Вручную поднимаете темы через 🚀 (первый раз)
3. Бот ждет указанный интервал (например, 12 часов) 3. Бот ждёт указанный интервал (например, 12 часов)
4. Автоматически поднимает темы, у которых прошло 12+ часов 4. Автоматически поднимает темы, у которых прошёл интервал
5. Повторяет каждые 12 часов 5. Цикл повторяется
**Пример:**
``` ```
00:00 - Запуск бота 00:00 Запуск бота
00:05 - Вы вручную нажали "Поднять темы" 00:05 — Ручное поднятие
12:05 - Автоматическое поднятие 12:05 Автоматическое поднятие
24:05 - Автоматическое поднятие 24:05 Автоматическое поднятие
36:05 - Автоматическое поднятие
``` ```
## 🏗️ Архитектура ## 🏗️ Архитектура
### Модульная структура
``` ```
app.py # Main bot logic with FSM app.py # Бот: хендлеры, middleware аутентификации, авто-бамп
config_manager.py # Type-safe configuration config_manager.py # Загрузка и валидация .env
api_client.py # Lolz batch API client with retry logic api_client.py # Lolz batch API клиент с retry логикой
database.py # SQLite with connection pooling database.py # SQLite: темы, настройки, статистика
``` ```
### Ключевые улучшения ### Что внутри БД (SQLite через aiosqlite)
- **Batch API**: До 10 тем за один запрос (10x эффективнее) | Таблица | Назначение |
- **Type Safety**: Полная типизация с Python 3.12+ (PEP 695) |---------|-----------|
- **SOLID Principles**: Каждый модуль имеет одну ответственность | `threads` | Темы: ID, название, дата последнего бампа |
- **Connection Pooling**: Эффективное управление соединениями | `settings` | Динамические настройки (интервал, batch_size, автобамп) |
- **Error Handling**: Специфичные исключения для каждого случая | `bump_stats` | Персистентная статистика (всего/успешно/последний бамп) |
- **Async Context Managers**: Автоматическая очистка ресурсов
- **Immutable Config**: Frozen dataclasses для безопасности
- **Logging**: Структурированные логи с rotation
- **Graceful Shutdown**: Корректное завершение всех задач
### Batch API Implementation ### Ключевые особенности
**Как работает:** - **Batch API** — до 10 тем за запрос (экономия лимитов в 10 раз)
```python - **Динамические настройки** — меняются в БД, подхватываются ботом мгновенно
# Старый способ (10 запросов): - **Персистентная статистика** — не теряется при перезапуске
for thread_id in [1, 2, 3, 4, 5, 6, 7, 8, 9, 10]: - **Аутентификация** — middleware проверяет Telegram ID админа
POST /threads/{thread_id}/bump - **Async context managers** — автоматическая очистка ресурсов
- **Frozen dataclasses с `__slots__`** — иммутабельность и экономия памяти
# Новый способ (1 batch запрос): - **WAL-режим SQLite** — лучшая конкурентность
POST /batch
[
{"uri": "/threads/1/bump", "method": "POST"},
{"uri": "/threads/2/bump", "method": "POST"},
...
{"uri": "/threads/10/bump", "method": "POST"}
]
```
**Преимущества:**
- Экономия API лимитов в 10 раз
- Скорость выполнения увеличена в ~10 раз
- Меньше нагрузка на сеть
- Атомарная обработка результатов
### Edge Cases
- Rate limiting от API (обработка в batch ответах)
- Network timeouts (retry для batch запросов)
- Invalid tokens (валидация перед batch)
- Partial batch failures (индивидуальная обработка каждого результата)
- Concurrent access (connection pooling)
- Database locks (async context managers)
- Memory leaks prevention (frozen dataclasses, proper cleanup)
## 📝 Логи ## 📝 Логи
- **Консоль** - основные события - **Консоль** основные события
- **bot.log** - детальные логи с traceback - **bot.log** детальные логи с traceback
## 🔧 Технические детали ## 🔧 Зависимости
- **Python**: 3.12+ - **Python**: 3.12+
- **aiogram**: 3.4.1 - **aiogram**: 3.15.0
- **aiohttp**: 3.10.0+ - **aiohttp**: bundled with aiogram
- **aiosqlite**: 0.20.0+ - **aiosqlite**: 0.20.0+
- **База данных**: SQLite с индексами - **python-dotenv**: 1.0.0+
- **API**: https://prod-api.lolz.live
- **Batch API**: До 10 запросов за 1 batch (каждый = 0.1 batch)
## 🐛 Troubleshooting ## 🐛 Troubleshooting
### Ошибка: "bot.api_token not configured" ### Ошибка: «BOT_API_TOKEN not configured»
Заполните `config.json` реальными токенами. Проверьте `.env` — заполните `BOT_API_TOKEN` реальным токеном от @BotFather.
### Ошибка: "Invalid API token" ### Ошибка: «ADMIN_USER_ID not configured»
Укажите в `.env` ваш Telegram ID (получить через @userinfobot).
### Ошибка: «API_AUTH_TOKEN not configured»
Проверьте токен на [zelenka.guru/account/api](https://zelenka.guru/account/api). Проверьте токен на [zelenka.guru/account/api](https://zelenka.guru/account/api).
### Ошибка: "Network error" ### Ошибка: «⛔ Доступ запрещён»
Проверьте интернет-соединение и доступность API. Ботом может управлять только пользователь, чей Telegram ID указан в `ADMIN_USER_ID`.
### Темы не поднимаются автоматически ### Темы не поднимаются автоматически
1. Проверьте `enable_auto_bump: true` в config.json 1. Проверьте, что автобамп включён (🔄 Toggle Auto-Bump)
2. Сделайте первое поднятие вручную через 🚀 2. Сделайте первое поднятие вручную через 🚀
3. Проверьте логи в bot.log 3. Проверьте логи в `bot.log`
## 👤 Автор ## 👤 Автор
[QIYANA](https://lolz.live/kqlol/) - создатель бота [QIYANA](https://lolz.live/kqlol/)
--- ---
+40 -7
View File
@@ -352,10 +352,10 @@ class APIClient:
continue continue
job_data = jobs[uri] job_data = jobs[uri]
logger.debug(f"Thread {thread_id} bump response: {job_data}") logger.info(f"Thread {thread_id} raw bump response: {job_data}")
# Empty dict means success for bump # Empty list [] or empty dict {} means success
if isinstance(job_data, dict) and len(job_data) == 0: if isinstance(job_data, (list, dict)) and len(job_data) == 0:
logger.info(f"Thread {thread_id} bumped successfully (empty response)") logger.info(f"Thread {thread_id} bumped successfully (empty response)")
results.append(BumpResult( results.append(BumpResult(
success=True, success=True,
@@ -364,7 +364,6 @@ class APIClient:
status=BumpStatus.SUCCESS status=BumpStatus.SUCCESS
)) ))
elif isinstance(job_data, dict) and "errors" in job_data: elif isinstance(job_data, dict) and "errors" in job_data:
# Has errors
error_msg = self._extract_error_message(str(job_data["errors"])) error_msg = self._extract_error_message(str(job_data["errors"]))
logger.error( logger.error(
f"Thread {thread_id} bump failed | " f"Thread {thread_id} bump failed | "
@@ -377,12 +376,46 @@ class APIClient:
thread_id=thread_id, thread_id=thread_id,
status=BumpStatus.ERROR status=BumpStatus.ERROR
)) ))
elif isinstance(job_data, dict) and "_job_result" in job_data:
job_result = str(job_data.get("_job_result", ""))
job_message = str(job_data.get("_job_message", ""))
if job_result == "error":
error_text = job_message
if not error_text.strip():
errors = job_data.get("errors")
if errors:
if isinstance(errors, list):
error_text = str(errors[0]) if errors else ""
elif isinstance(errors, str):
error_text = errors
else:
error_text = str(errors)
if not error_text.strip():
error_text = str(job_data.get("error", ""))
error_msg = self._extract_error_message(error_text) or "Ошибка API (см. логи)"
logger.error(
f"Thread {thread_id} bump failed | "
f"Job error: {error_msg}"
)
results.append(BumpResult(
success=False,
message=f"Тема {thread_id}: {error_msg}",
thread_id=thread_id,
status=BumpStatus.ERROR
))
else:
logger.info(f"Thread {thread_id} bumped successfully (job_result={job_result}, job_message={job_message})")
results.append(BumpResult(
success=True,
message=f"✅ Тема {thread_id} поднята успешно",
thread_id=thread_id,
status=BumpStatus.SUCCESS
))
else: else:
# Unknown response logger.warning(f"Thread {thread_id} unknown bump response (type={type(job_data).__name__}): {job_data}")
logger.warning(f"Thread {thread_id} unknown bump response: {job_data}")
results.append(BumpResult( results.append(BumpResult(
success=False, success=False,
message=f"Тема {thread_id}: Неизвестный ответ", message=f"Тема {thread_id}: Неизвестный ответ ({type(job_data).__name__})",
thread_id=thread_id, thread_id=thread_id,
status=BumpStatus.ERROR status=BumpStatus.ERROR
)) ))
+822 -661
View File
File diff suppressed because it is too large Load Diff
+108 -120
View File
@@ -1,120 +1,108 @@
"""Configuration management with validation and type safety.""" """Configuration management with .env validation and type safety."""
import json import os
from pathlib import Path from dataclasses import dataclass
from dataclasses import dataclass from typing import Self
from typing import Self
from dotenv import load_dotenv
@dataclass(frozen=True, slots=True)
class BotConfig: @dataclass(frozen=True, slots=True)
"""Telegram bot configuration.""" class BotConfig:
api_token: str api_token: str
img_url: str img_url: str
author_url: str author_url: str
@dataclass(frozen=True, slots=True) @dataclass(frozen=True, slots=True)
class APIConfig: class APIConfig:
"""Lolz API configuration.""" base_url: str
base_url: str auth_token: str
auth_token: str batch_size: int
batch_size: int
@dataclass(frozen=True, slots=True)
@dataclass(frozen=True, slots=True) class DatabaseConfig:
class DatabaseConfig: path: str
"""Database configuration."""
path: str
@dataclass(frozen=True, slots=True)
class SchedulingConfig:
@dataclass(frozen=True, slots=True) bump_interval_hours: float
class SchedulingConfig: bump_delay_seconds: float
"""Scheduling configuration.""" enable_auto_bump: bool
bump_interval_hours: float
bump_delay_seconds: float
enable_auto_bump: bool @dataclass(frozen=True, slots=True)
class Config:
bot: BotConfig
@dataclass(frozen=True, slots=True) api: APIConfig
class Config: database: DatabaseConfig
"""Application configuration.""" scheduling: SchedulingConfig
bot: BotConfig admin_user_id: int
api: APIConfig
database: DatabaseConfig @classmethod
scheduling: SchedulingConfig def load(cls, env_path: str = ".env") -> Self:
if not os.path.exists(env_path):
@classmethod raise FileNotFoundError(f".env file not found: {env_path}")
def load(cls, config_path: str = "config.json") -> Self:
"""Load and validate configuration from JSON file.""" load_dotenv(env_path)
path = Path(config_path)
if not path.exists(): bot_token = os.getenv("BOT_API_TOKEN", "")
raise FileNotFoundError(f"Config file not found: {config_path}") api_token = os.getenv("API_AUTH_TOKEN", "")
admin_user_id = int(os.getenv("ADMIN_USER_ID", "0"))
with open(path, encoding="utf-8") as f:
data = json.load(f) cls._validate_tokens(bot_token, api_token, admin_user_id)
# Validate required fields bot = BotConfig(
cls._validate_config(data) api_token=bot_token,
img_url=os.getenv("BOT_IMG_URL", ""),
return cls( author_url=os.getenv("BOT_AUTHOR_URL", ""),
bot=BotConfig( )
api_token=data["bot"]["api_token"],
img_url=data["bot"]["img_url"], batch_size = int(os.getenv("API_BATCH_SIZE", "10"))
author_url=data["bot"]["author_url"] if batch_size < 1 or batch_size > 10:
), raise ValueError("API_BATCH_SIZE must be between 1 and 10")
api=APIConfig(
base_url=data["api"]["base_url"].rstrip("/"), api = APIConfig(
auth_token=data["api"]["auth_token"], base_url=os.getenv("API_BASE_URL", "").rstrip("/"),
batch_size=int(data["api"].get("batch_size", 10)) auth_token=api_token,
), batch_size=batch_size,
database=DatabaseConfig( )
path=data["database"]["path"]
), db = DatabaseConfig(
scheduling=SchedulingConfig( path=os.getenv("DB_PATH", "threads.db"),
bump_interval_hours=float(data["scheduling"]["bump_interval_hours"]), )
bump_delay_seconds=float(data["scheduling"].get("bump_delay_seconds", 2)),
enable_auto_bump=bool(data["scheduling"]["enable_auto_bump"]) interval = float(os.getenv("BUMP_INTERVAL_HOURS", "12"))
) if interval <= 0:
) raise ValueError("BUMP_INTERVAL_HOURS must be positive")
@staticmethod scheduling = SchedulingConfig(
def _validate_config(data: dict) -> None: bump_interval_hours=interval,
"""Validate configuration structure and required fields.""" bump_delay_seconds=float(os.getenv("BUMP_DELAY_SECONDS", "2")),
required_fields = { enable_auto_bump=os.getenv("ENABLE_AUTO_BUMP", "true").lower() == "true",
"bot": ["api_token", "img_url", "author_url"], )
"api": ["base_url", "auth_token"],
"database": ["path"], return cls(
"scheduling": ["bump_interval_hours", "enable_auto_bump"] bot=bot,
} api=api,
database=db,
for section, fields in required_fields.items(): scheduling=scheduling,
if section not in data: admin_user_id=admin_user_id,
raise ValueError(f"Missing config section: {section}") )
for field in fields: @staticmethod
if field not in data[section]: def _validate_tokens(bot_token: str, api_token: str, admin_user_id: int) -> None:
raise ValueError(f"Missing config field: {section}.{field}") if not bot_token or "YOUR_" in bot_token:
raise ValueError(
value = data[section][field] "BOT_API_TOKEN not configured — set your Telegram bot token in .env"
if isinstance(value, str) and not value.strip(): )
raise ValueError(f"Empty config field: {section}.{field}") if not api_token or "YOUR_" in api_token:
raise ValueError(
# Validate token formats "API_AUTH_TOKEN not configured — set your Lolz API token in .env"
bot_token = data["bot"]["api_token"] )
if "YOUR_" in bot_token or not bot_token: if admin_user_id == 0:
raise ValueError("bot.api_token not configured - please set your Telegram bot token") raise ValueError(
"ADMIN_USER_ID not configured — set your Telegram user ID in .env"
api_token = data["api"]["auth_token"] )
if "YOUR_" in api_token or not api_token:
raise ValueError("api.auth_token not configured - please set your Lolz API token")
# Validate interval
interval = data["scheduling"]["bump_interval_hours"]
if not isinstance(interval, (int, float)) or interval <= 0:
raise ValueError("scheduling.bump_interval_hours must be positive number")
# Validate batch size
batch_size = data["api"].get("batch_size", 10)
if not isinstance(batch_size, int) or batch_size < 1 or batch_size > 10:
raise ValueError("api.batch_size must be between 1 and 10")
+295 -259
View File
@@ -1,259 +1,295 @@
"""Database layer with connection pooling and type safety.""" """Database layer with dynamic settings and stats persistence."""
import aiosqlite import aiosqlite
import logging import logging
from typing import Self from typing import Self
from dataclasses import dataclass from dataclasses import dataclass
from config_manager import Config
logger = logging.getLogger(__name__)
logger = logging.getLogger(__name__)
@dataclass(frozen=True, slots=True)
class Thread:
"""Thread model with immutable fields.""" @dataclass(frozen=True, slots=True)
id: str class Thread:
title: str id: str
last_bumped: str | None = None title: str
last_bumped: str | None = None
class Database:
"""SQLite database manager with connection pooling and proper error handling.""" @dataclass(frozen=True, slots=True)
class BumpStats:
__slots__ = ("_db_path", "_connection") total_bumps: int
successful_bumps: int
def __init__(self, db_path: str) -> None: last_bump_time: str | None = None
if not db_path:
raise ValueError("db_path is required")
class Database:
self._db_path = db_path __slots__ = ("_db_path", "_connection")
self._connection: aiosqlite.Connection | None = None
def __init__(self, db_path: str) -> None:
async def __aenter__(self) -> Self: if not db_path:
"""Async context manager entry.""" raise ValueError("db_path is required")
await self.connect()
return self self._db_path = db_path
self._connection: aiosqlite.Connection | None = None
async def __aexit__(self, exc_type, exc_val, exc_tb) -> None:
"""Async context manager exit.""" async def __aenter__(self) -> Self:
await self.close() await self.connect()
return self
async def connect(self) -> None:
"""Establish database connection and initialize schema.""" async def __aexit__(self, exc_type, exc_val, exc_tb) -> None:
if self._connection is None: await self.close()
try:
self._connection = await aiosqlite.connect(self._db_path) async def connect(self) -> None:
# Enable WAL mode for better concurrency if self._connection is None:
await self._connection.execute("PRAGMA journal_mode=WAL") try:
await self._initialize_schema() self._connection = await aiosqlite.connect(self._db_path)
logger.info(f"Database connected: {self._db_path}") await self._connection.execute("PRAGMA journal_mode=WAL")
except Exception as e: await self._connection.execute("PRAGMA foreign_keys=ON")
logger.error(f"Failed to connect to database: {e}") await self._initialize_schema()
raise logger.info(f"Database connected: {self._db_path}")
except Exception as e:
async def close(self) -> None: logger.error(f"Failed to connect to database: {e}")
"""Close database connection safely.""" raise
if self._connection:
try: async def close(self) -> None:
await self._connection.close() if self._connection:
self._connection = None try:
logger.info("Database connection closed") await self._connection.close()
except Exception as e: self._connection = None
logger.error(f"Error closing database connection: {e}") logger.info("Database connection closed")
except Exception as e:
async def _initialize_schema(self) -> None: logger.error(f"Error closing database connection: {e}")
"""Initialize database schema with tables and indexes."""
if not self._connection: async def _initialize_schema(self) -> None:
raise RuntimeError("Database not connected") if not self._connection:
raise RuntimeError("Database not connected")
try:
await self._connection.execute(""" try:
CREATE TABLE IF NOT EXISTS threads ( await self._connection.execute("""
id TEXT PRIMARY KEY, CREATE TABLE IF NOT EXISTS threads (
title TEXT NOT NULL, id TEXT PRIMARY KEY,
last_bumped TEXT, title TEXT NOT NULL,
created_at TEXT DEFAULT CURRENT_TIMESTAMP last_bumped TEXT,
) created_at TEXT DEFAULT CURRENT_TIMESTAMP
""") )
""")
# Create index for efficient queries on last_bumped
await self._connection.execute(""" await self._connection.execute("""
CREATE INDEX IF NOT EXISTS idx_last_bumped CREATE INDEX IF NOT EXISTS idx_last_bumped
ON threads(last_bumped) ON threads(last_bumped)
""") """)
await self._connection.commit() await self._connection.execute("""
logger.info("Database schema initialized") CREATE TABLE IF NOT EXISTS settings (
except Exception as e: key TEXT PRIMARY KEY,
logger.error(f"Failed to initialize schema: {e}") value TEXT NOT NULL,
raise updated_at TEXT DEFAULT CURRENT_TIMESTAMP
)
def _ensure_connected(self) -> None: """)
"""Ensure database is connected before operations."""
if not self._connection: await self._connection.execute("""
raise RuntimeError("Database not connected. Call connect() first.") CREATE TABLE IF NOT EXISTS bump_stats (
id INTEGER PRIMARY KEY CHECK (id = 1),
async def add_thread(self, thread_id: str, title: str) -> bool: total_bumps INTEGER DEFAULT 0,
""" successful_bumps INTEGER DEFAULT 0,
Add thread to database. last_bump_time TEXT
)
Args: """)
thread_id: Unique thread identifier
title: Thread title await self._connection.commit()
logger.info("Database schema initialized")
Returns: except Exception as e:
True if added successfully, False if thread already exists logger.error(f"Failed to initialize schema: {e}")
""" raise
self._ensure_connected()
async def seed_defaults(self, config: Config) -> None:
try: await self._seed_settings(config)
await self._connection.execute( await self._seed_stats()
"INSERT INTO threads (id, title) VALUES (?, ?)",
(thread_id, title) async def _seed_settings(self, config: Config) -> None:
) defaults = {
await self._connection.commit() "bump_interval_hours": str(config.scheduling.bump_interval_hours),
return True "bump_delay_seconds": str(config.scheduling.bump_delay_seconds),
except aiosqlite.IntegrityError: "enable_auto_bump": str(config.scheduling.enable_auto_bump).lower(),
# Thread already exists (PRIMARY KEY constraint) "batch_size": str(config.api.batch_size),
return False }
except Exception as e: for key, value in defaults.items():
logger.error(f"Error adding thread {thread_id}: {e}") await self._connection.execute(
raise "INSERT OR IGNORE INTO settings (key, value) VALUES (?, ?)",
(key, value),
async def get_all_threads(self) -> list[Thread]: )
""" await self._connection.commit()
Get all threads ordered by ID.
async def _seed_stats(self) -> None:
Returns: await self._connection.execute(
List of Thread objects "INSERT OR IGNORE INTO bump_stats (id, total_bumps, successful_bumps) VALUES (1, 0, 0)"
""" )
self._ensure_connected() await self._connection.commit()
try: def _ensure_connected(self) -> None:
async with self._connection.execute( if not self._connection:
"SELECT id, title, last_bumped FROM threads ORDER BY id" raise RuntimeError("Database not connected. Call connect() first.")
) as cursor:
rows = await cursor.fetchall() # ─── Settings ───────────────────────────────────────────────
return [
Thread(id=row[0], title=row[1], last_bumped=row[2]) async def get_setting(self, key: str, default: str | None = None) -> str | None:
for row in rows self._ensure_connected()
] async with self._connection.execute(
except Exception as e: "SELECT value FROM settings WHERE key = ?", (key,)
logger.error(f"Error fetching all threads: {e}") ) as cursor:
raise row = await cursor.fetchone()
return row[0] if row else default
async def get_threads_to_bump(self, interval_hours: float) -> list[Thread]:
""" async def set_setting(self, key: str, value: str) -> None:
Get threads ready for bumping based on interval. self._ensure_connected()
await self._connection.execute(
Args: "INSERT OR REPLACE INTO settings (key, value, updated_at) VALUES (?, ?, datetime('now'))",
interval_hours: Minimum hours since last bump (key, value),
)
Returns: await self._connection.commit()
List of Thread objects ready to bump
""" async def get_all_settings(self) -> dict[str, str]:
self._ensure_connected() self._ensure_connected()
async with self._connection.execute("SELECT key, value FROM settings") as cursor:
if interval_hours < 0: rows = await cursor.fetchall()
raise ValueError("interval_hours must be non-negative") return {row[0]: row[1] for row in rows}
try: # ─── Bump Statistics ────────────────────────────────────────
async with self._connection.execute(
""" async def get_bump_stats(self) -> BumpStats:
SELECT id, title, last_bumped FROM threads self._ensure_connected()
WHERE last_bumped IS NULL async with self._connection.execute(
OR datetime(last_bumped, '+' || ? || ' hours') <= datetime('now') "SELECT total_bumps, successful_bumps, last_bump_time FROM bump_stats WHERE id = 1"
ORDER BY last_bumped ASC NULLS FIRST ) as cursor:
""", row = await cursor.fetchone()
(interval_hours,) if row:
) as cursor: return BumpStats(
rows = await cursor.fetchall() total_bumps=row[0],
return [ successful_bumps=row[1],
Thread(id=row[0], title=row[1], last_bumped=row[2]) last_bump_time=row[2],
for row in rows )
] return BumpStats(total_bumps=0, successful_bumps=0)
except Exception as e:
logger.error(f"Error fetching threads to bump: {e}") async def increment_bump_stats(self, success_count: int, total_count: int) -> None:
raise self._ensure_connected()
await self._connection.execute(
async def update_last_bumped(self, thread_id: str) -> None: """
""" UPDATE bump_stats SET
Update last bumped timestamp for thread to current time. total_bumps = total_bumps + ?,
successful_bumps = successful_bumps + ?,
Args: last_bump_time = datetime('now')
thread_id: Thread identifier to update WHERE id = 1
""" """,
self._ensure_connected() (total_count, success_count),
)
try: await self._connection.commit()
await self._connection.execute(
"UPDATE threads SET last_bumped = datetime('now') WHERE id = ?", # ─── Threads ────────────────────────────────────────────────
(thread_id,)
) async def add_thread(self, thread_id: str, title: str) -> bool:
await self._connection.commit() self._ensure_connected()
except Exception as e: try:
logger.error(f"Error updating last_bumped for thread {thread_id}: {e}") await self._connection.execute(
raise "INSERT INTO threads (id, title) VALUES (?, ?)",
(thread_id, title),
async def delete_thread(self, thread_id: str) -> bool: )
""" await self._connection.commit()
Delete thread from database. return True
except aiosqlite.IntegrityError:
Args: return False
thread_id: Thread identifier to delete except Exception as e:
logger.error(f"Error adding thread {thread_id}: {e}")
Returns: raise
True if deleted successfully, False if thread not found
""" async def update_thread_title(self, thread_id: str, title: str) -> bool:
self._ensure_connected() self._ensure_connected()
try:
try: cursor = await self._connection.execute(
cursor = await self._connection.execute( "UPDATE threads SET title = ? WHERE id = ?", (title, thread_id)
"DELETE FROM threads WHERE id = ?", )
(thread_id,) await self._connection.commit()
) return cursor.rowcount > 0
await self._connection.commit() except Exception as e:
return cursor.rowcount > 0 logger.error(f"Error updating title for thread {thread_id}: {e}")
except Exception as e: raise
logger.error(f"Error deleting thread {thread_id}: {e}")
raise async def get_all_threads(self) -> list[Thread]:
self._ensure_connected()
async def get_thread_count(self) -> int: try:
""" async with self._connection.execute(
Get total number of threads in database. "SELECT id, title, last_bumped FROM threads ORDER BY id"
) as cursor:
Returns: rows = await cursor.fetchall()
Total thread count return [Thread(id=row[0], title=row[1], last_bumped=row[2]) for row in rows]
""" except Exception as e:
self._ensure_connected() logger.error(f"Error fetching all threads: {e}")
raise
try:
async with self._connection.execute("SELECT COUNT(*) FROM threads") as cursor: async def get_threads_to_bump(self, interval_hours: float) -> list[Thread]:
row = await cursor.fetchone() self._ensure_connected()
return row[0] if row else 0 if interval_hours < 0:
except Exception as e: raise ValueError("interval_hours must be non-negative")
logger.error(f"Error getting thread count: {e}") try:
raise async with self._connection.execute(
"""
async def thread_exists(self, thread_id: str) -> bool: SELECT id, title, last_bumped FROM threads
""" WHERE last_bumped IS NULL
Check if thread exists in database. OR datetime(last_bumped, '+' || ? || ' hours') <= datetime('now')
ORDER BY last_bumped ASC NULLS FIRST
Args: """,
thread_id: Thread identifier to check (interval_hours,),
) as cursor:
Returns: rows = await cursor.fetchall()
True if thread exists, False otherwise return [Thread(id=row[0], title=row[1], last_bumped=row[2]) for row in rows]
""" except Exception as e:
self._ensure_connected() logger.error(f"Error fetching threads to bump: {e}")
raise
try:
async with self._connection.execute( async def update_last_bumped(self, thread_id: str) -> None:
"SELECT 1 FROM threads WHERE id = ? LIMIT 1", self._ensure_connected()
(thread_id,) try:
) as cursor: await self._connection.execute(
row = await cursor.fetchone() "UPDATE threads SET last_bumped = datetime('now') WHERE id = ?",
return row is not None (thread_id,),
except Exception as e: )
logger.error(f"Error checking thread existence {thread_id}: {e}") await self._connection.commit()
raise except Exception as e:
logger.error(f"Error updating last_bumped for thread {thread_id}: {e}")
raise
async def delete_thread(self, thread_id: str) -> bool:
self._ensure_connected()
try:
cursor = await self._connection.execute(
"DELETE FROM threads WHERE id = ?", (thread_id,)
)
await self._connection.commit()
return cursor.rowcount > 0
except Exception as e:
logger.error(f"Error deleting thread {thread_id}: {e}")
raise
async def get_thread_count(self) -> int:
self._ensure_connected()
try:
async with self._connection.execute("SELECT COUNT(*) FROM threads") as cursor:
row = await cursor.fetchone()
return row[0] if row else 0
except Exception as e:
logger.error(f"Error getting thread count: {e}")
raise
async def thread_exists(self, thread_id: str) -> bool:
self._ensure_connected()
try:
async with self._connection.execute(
"SELECT 1 FROM threads WHERE id = ? LIMIT 1", (thread_id,)
) as cursor:
row = await cursor.fetchone()
return row is not None
except Exception as e:
logger.error(f"Error checking thread existence {thread_id}: {e}")
raise
+3 -2
View File
@@ -1,2 +1,3 @@
aiogram==3.15.0 aiogram==3.15.0
aiosqlite>=0.20.0 aiosqlite>=0.20.0
python-dotenv>=1.0.0