forked from FOSS/ASF-Control-Bot
353 lines
13 KiB
Python
353 lines
13 KiB
Python
import asyncio
|
|
import hashlib
|
|
import io
|
|
import json
|
|
import zipfile
|
|
from dataclasses import dataclass
|
|
from datetime import datetime, timezone
|
|
from pathlib import PurePosixPath
|
|
|
|
from aiogram import Router
|
|
from aiogram.fsm.context import FSMContext
|
|
from aiogram.types import CallbackQuery, Message
|
|
|
|
from bot.api.asf import get_bot, save_bot_config
|
|
from bot.config import ADMIN_ID
|
|
from bot.handlers.bots import send_bots_menu
|
|
from bot.logging_utils import get_logger
|
|
from bot.services.balance import (
|
|
InsufficientBalance,
|
|
increase_balance,
|
|
debit_user,
|
|
get_bot_slot_price,
|
|
get_user_balance,
|
|
)
|
|
from bot.services.bot_accounts import (
|
|
bot_name_exists,
|
|
create_bot_account,
|
|
user_hash_exists,
|
|
)
|
|
from bot.states import UploadFlow
|
|
from bot.ui.keyboards import back_keyboard
|
|
|
|
router = Router()
|
|
logger = get_logger(__name__)
|
|
_groups: dict[str, list[Message]] = {}
|
|
_group_tasks: dict[str, asyncio.Task] = {}
|
|
|
|
|
|
@router.callback_query(lambda c: c.data == "upload")
|
|
async def upload_menu(callback: CallbackQuery, state: FSMContext) -> None:
|
|
await callback.answer()
|
|
await state.set_state(UploadFlow.waiting_upload)
|
|
|
|
balance = await get_user_balance(callback.from_user.id)
|
|
slot_price = await get_bot_slot_price()
|
|
|
|
await callback.message.edit_text(
|
|
"Отправьте документы .json или .zip с JSON-конфигурациями.\n\n"
|
|
f"💰 Ваш баланс: {balance} ₽\n"
|
|
f"💳 Цена одного слота: {slot_price} ₽",
|
|
reply_markup=back_keyboard(callback_data="bots"),
|
|
)
|
|
|
|
|
|
def _decode_config(raw: bytes) -> dict:
|
|
value = json.loads(raw.decode("utf-8-sig"))
|
|
if not isinstance(value, dict):
|
|
raise ValueError("JSON должен содержать объект конфигурации")
|
|
name = value.get("SteamLogin")
|
|
if not isinstance(name, str) or not name.strip():
|
|
raise ValueError("отсутствует непустой SteamLogin")
|
|
return value
|
|
|
|
|
|
def _zip_entries(raw: bytes) -> list[tuple[str, bytes]]:
|
|
entries = []
|
|
with zipfile.ZipFile(io.BytesIO(raw)) as archive:
|
|
for item in archive.infolist():
|
|
name = item.filename.replace("\\", "/")
|
|
path = PurePosixPath(name)
|
|
if path.is_absolute() or ".." in path.parts:
|
|
raise ValueError("ZIP содержит небезопасный путь")
|
|
if item.is_dir():
|
|
continue
|
|
if not name.lower().endswith(".json"):
|
|
raise ValueError(f"ZIP содержит неподдерживаемый файл: {name}")
|
|
entries.append((name, archive.read(item)))
|
|
return entries
|
|
|
|
|
|
async def _download(message: Message) -> bytes:
|
|
buffer = io.BytesIO()
|
|
await message.bot.download(message.document.file_id, destination=buffer)
|
|
return buffer.getvalue()
|
|
|
|
|
|
@dataclass
|
|
class _Candidate:
|
|
filename: str
|
|
config: dict
|
|
digest: str
|
|
|
|
|
|
def _digest(config: dict) -> str:
|
|
canonical = json.dumps(
|
|
config, ensure_ascii=False, sort_keys=True, separators=(",", ":")
|
|
)
|
|
return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
|
|
|
|
|
|
async def _preflight(
|
|
user_id: int, files: list[tuple[str, bytes]]
|
|
) -> tuple[list[_Candidate], list[str]]:
|
|
candidates: list[_Candidate] = []
|
|
errors: list[str] = []
|
|
parsed: list[tuple[str, dict, str]] = []
|
|
|
|
for filename, raw in files:
|
|
try:
|
|
config = _decode_config(raw)
|
|
name = config["SteamLogin"].strip()
|
|
digest = _digest(config)
|
|
parsed.append((filename, config, digest))
|
|
|
|
except (ValueError, json.JSONDecodeError, UnicodeDecodeError) as error:
|
|
logger.warning(
|
|
"Предварительная проверка отклонена: пользователь=%s файл=%s ошибка=%s",
|
|
user_id,
|
|
filename,
|
|
error,
|
|
)
|
|
errors.append(f"❌ {filename}: {error}")
|
|
|
|
except Exception as error:
|
|
logger.exception(
|
|
"Ошибка предварительной проверки: пользователь=%s файл=%s",
|
|
user_id,
|
|
filename,
|
|
)
|
|
errors.append(f"❌ {filename}: ошибка проверки ({error})")
|
|
|
|
name_counts: dict[str, int] = {}
|
|
hash_counts: dict[str, int] = {}
|
|
for _, config, digest in parsed:
|
|
name_counts[config["SteamLogin"].strip()] = (
|
|
name_counts.get(config["SteamLogin"].strip(), 0) + 1
|
|
)
|
|
hash_counts[digest] = hash_counts.get(digest, 0) + 1
|
|
|
|
for filename, config, digest in parsed:
|
|
name = config["SteamLogin"].strip()
|
|
try:
|
|
if name_counts[name] > 1:
|
|
raise ValueError("дублируется имя бота в этой загрузке")
|
|
|
|
if hash_counts[digest] > 1 or await user_hash_exists(user_id, digest):
|
|
raise ValueError(
|
|
"такая конфигурация уже загружалась или дублируется в загрузке"
|
|
)
|
|
|
|
if await bot_name_exists(name):
|
|
raise ValueError("бот с таким именем уже есть в базе данных")
|
|
existing = get_bot(name)
|
|
|
|
if existing.status_code == 200 and existing.json().get("Result"):
|
|
raise ValueError("бот с таким именем уже есть в ASF")
|
|
|
|
candidates.append(_Candidate(filename, config, digest))
|
|
logger.info(
|
|
"Предварительная проверка пройдена: пользователь=%s файл=%s бот=%s",
|
|
user_id,
|
|
filename,
|
|
name,
|
|
)
|
|
|
|
except (ValueError, json.JSONDecodeError, UnicodeDecodeError) as error:
|
|
logger.warning(
|
|
"Предварительная проверка отклонена: пользователь=%s файл=%s ошибка=%s",
|
|
user_id,
|
|
filename,
|
|
error,
|
|
)
|
|
errors.append(f"❌ {filename}: {error}")
|
|
|
|
except Exception as error:
|
|
logger.exception(
|
|
"Ошибка предварительной проверки: пользователь=%s файл=%s",
|
|
user_id,
|
|
filename,
|
|
)
|
|
errors.append(f"❌ {filename}: ошибка проверки ({error})")
|
|
|
|
return candidates, errors
|
|
|
|
|
|
async def _process(user_id: int, files: list[tuple[str, bytes]]) -> str:
|
|
candidates, results = await _preflight(user_id, files)
|
|
if not candidates:
|
|
return "\n".join(results) or "❌ Нет конфигураций для загрузки"
|
|
|
|
is_admin = user_id == ADMIN_ID
|
|
slot_price = await get_bot_slot_price()
|
|
total = slot_price * len(candidates)
|
|
reserved = False
|
|
if not is_admin:
|
|
try:
|
|
await debit_user(user_id, total)
|
|
reserved = True
|
|
logger.info(
|
|
"Средства зарезервированы: пользователь=%s сумма=%s ₽ количество=%s",
|
|
user_id,
|
|
total,
|
|
len(candidates),
|
|
)
|
|
|
|
except InsufficientBalance:
|
|
balance = await get_user_balance(user_id)
|
|
logger.warning(
|
|
"Недостаточно средств для загрузки: пользователь=%s баланс=%s ₽ требуется=%s ₽",
|
|
user_id,
|
|
balance,
|
|
total,
|
|
)
|
|
|
|
results.append(
|
|
f"❌ Недостаточно средств: нужно {total} ₽ для {len(candidates)} конф. "
|
|
f"(доступно {balance} ₽)\n"
|
|
)
|
|
|
|
results.extend(
|
|
f"❌ {candidate.filename}: загрузка не начата"
|
|
for candidate in candidates
|
|
)
|
|
return "\n".join(results)
|
|
|
|
for candidate in candidates:
|
|
name = candidate.config["SteamLogin"].strip()
|
|
try:
|
|
response = save_bot_config(name, candidate.config)
|
|
logger.info(
|
|
"Ответ ASF на загрузку: пользователь=%s бот=%s статус=%s",
|
|
user_id,
|
|
name,
|
|
response.status_code,
|
|
)
|
|
|
|
payload = response.json()
|
|
if response.status_code != 200 or not payload.get("Success"):
|
|
detail = (
|
|
payload.get("Message")
|
|
or payload.get("Result")
|
|
or f"HTTP {response.status_code}"
|
|
)
|
|
raise ValueError(f"ASF отклонил конфигурацию: {detail}")
|
|
|
|
await create_bot_account(
|
|
user_id,
|
|
bot_name=name,
|
|
steam_id=(
|
|
str(candidate.config.get("SteamID"))
|
|
if candidate.config.get("SteamID")
|
|
else None
|
|
),
|
|
source_filename=candidate.filename[:255],
|
|
config_sha256=candidate.digest,
|
|
upload_state="attached",
|
|
is_attached=True,
|
|
attached_at=datetime.now(timezone.utc),
|
|
)
|
|
logger.info(
|
|
"Регистрация в базе завершена: пользователь=%s бот=%s", user_id, name
|
|
)
|
|
results.append(
|
|
f"✅ {candidate.filename}: ASF принял, аккаунт зарегистрирован ({name})"
|
|
)
|
|
|
|
except Exception as error:
|
|
logger.exception(
|
|
"Ошибка загрузки: пользователь=%s файл=%s бот=%s",
|
|
user_id,
|
|
candidate.filename,
|
|
name,
|
|
)
|
|
if reserved:
|
|
await increase_balance(user_id, slot_price)
|
|
logger.info(
|
|
"Средства возвращены: пользователь=%s сумма=%s ₽ файл=%s",
|
|
user_id,
|
|
slot_price,
|
|
candidate.filename,
|
|
)
|
|
results.append(f"❌ {candidate.filename}: {error}")
|
|
|
|
return "\n".join(results)
|
|
|
|
|
|
async def _finish_group(key: str, state: FSMContext) -> None:
|
|
await asyncio.sleep(0.8)
|
|
messages = _groups.pop(key, [])
|
|
_group_tasks.pop(key, None)
|
|
|
|
if not messages:
|
|
return
|
|
|
|
files = []
|
|
for message in messages:
|
|
filename = message.document.file_name or "config.json"
|
|
try:
|
|
files.append((filename, await _download(message)))
|
|
except Exception:
|
|
files.append((filename, b""))
|
|
|
|
await _send_result_and_menu(messages[0], messages[0].from_user.id, files)
|
|
await state.clear()
|
|
|
|
|
|
async def _send_result_and_menu(
|
|
message: Message, user_id: int, files: list[tuple[str, bytes]]
|
|
) -> None:
|
|
result = await _process(user_id, files)
|
|
await message.answer(result)
|
|
if "✅" in result:
|
|
await send_bots_menu(message, user_id)
|
|
|
|
|
|
@router.message(UploadFlow.waiting_upload, lambda message: bool(message.document))
|
|
async def upload_document(message: Message, state: FSMContext) -> None:
|
|
document = message.document
|
|
filename = document.file_name or "config.json"
|
|
|
|
if message.media_group_id:
|
|
key = f"{message.chat.id}:{message.media_group_id}"
|
|
_groups.setdefault(key, []).append(message)
|
|
if key not in _group_tasks:
|
|
_group_tasks[key] = asyncio.create_task(_finish_group(key, state))
|
|
return
|
|
|
|
try:
|
|
raw = await _download(message)
|
|
if filename.lower().endswith(".zip"):
|
|
files = _zip_entries(raw)
|
|
if not files:
|
|
await message.answer("❌ В архиве нет JSON-файлов")
|
|
else:
|
|
await _send_result_and_menu(message, message.from_user.id, files)
|
|
elif filename.lower().endswith(".json"):
|
|
await _send_result_and_menu(
|
|
message, message.from_user.id, [(filename, raw)]
|
|
)
|
|
else:
|
|
await message.answer("❌ Поддерживаются только документы .json и .zip")
|
|
except (zipfile.BadZipFile, ValueError) as error:
|
|
await message.answer(f"❌ {filename}: {error}")
|
|
finally:
|
|
await state.clear()
|
|
|
|
|
|
@router.message(UploadFlow.waiting_upload)
|
|
async def upload_non_document(message: Message, state: FSMContext) -> None:
|
|
await message.answer(
|
|
"❌ Отправьте конфигурацию именно документом .json или архивом .zip"
|
|
)
|