add restart storage

This commit is contained in:
Yuriy Yuriev
2026-05-23 00:07:44 +07:00
parent 68275386bf
commit 406a58c585
6 changed files with 384 additions and 18 deletions
+27
View File
@@ -1,8 +1,10 @@
import asyncio import asyncio
import hashlib import hashlib
import json
import secrets import secrets
import logging import logging
from datetime import datetime, timedelta from datetime import datetime, timedelta
from pathlib import Path
from typing import Dict, Optional from typing import Dict, Optional
from config.settings import settings from config.settings import settings
@@ -10,6 +12,7 @@ from config.settings import settings
logger = logging.getLogger(__name__) logger = logging.getLogger(__name__)
ADMIN_PASSWORD = settings.ADMIN_PASSWORD ADMIN_PASSWORD = settings.ADMIN_PASSWORD
_SESSIONS_FILE = Path("data/sessions.json")
class AuthManager: class AuthManager:
@@ -22,6 +25,7 @@ class AuthManager:
self._max_attempts = 5 self._max_attempts = 5
self._lockout_duration = 300 self._lockout_duration = 300
self._session_timeout = timedelta(minutes=session_timeout_minutes) self._session_timeout = timedelta(minutes=session_timeout_minutes)
self._restore_sessions()
@staticmethod @staticmethod
def _hash_password(password: str, salt: str = None) -> tuple: def _hash_password(password: str, salt: str = None) -> tuple:
@@ -53,14 +57,37 @@ class AuthManager:
self._sessions[user_id] = datetime.now() self._sessions[user_id] = datetime.now()
self._roles[user_id] = "admin" if is_admin else "user" self._roles[user_id] = "admin" if is_admin else "user"
self._failed_attempts.pop(user_id, None) self._failed_attempts.pop(user_id, None)
self._save_sessions()
logger.info(f"User {user_id} logged in as {self._roles[user_id]}") logger.info(f"User {user_id} logged in as {self._roles[user_id]}")
def logout(self, user_id: int): def logout(self, user_id: int):
self._sessions.pop(user_id, None) self._sessions.pop(user_id, None)
self._roles.pop(user_id, None) self._roles.pop(user_id, None)
self._failed_attempts.pop(user_id, None) self._failed_attempts.pop(user_id, None)
self._save_sessions()
logger.info(f"User {user_id} logged out") logger.info(f"User {user_id} logged out")
def _save_sessions(self) -> None:
try:
_SESSIONS_FILE.parent.mkdir(parents=True, exist_ok=True)
data = {str(uid): role for uid, role in self._roles.items()}
_SESSIONS_FILE.write_text(json.dumps(data), encoding="utf-8")
except Exception as e:
logger.warning(f"Failed to save sessions: {e}")
def _restore_sessions(self) -> None:
try:
if not _SESSIONS_FILE.exists():
return
data = json.loads(_SESSIONS_FILE.read_text(encoding="utf-8"))
for uid_str, role in data.items():
uid = int(uid_str)
self._sessions[uid] = datetime.now()
self._roles[uid] = role
logger.info(f"Restored {len(data)} sessions")
except Exception as e:
logger.warning(f"Failed to restore sessions: {e}")
def is_locked_out(self, user_id: int) -> bool: def is_locked_out(self, user_id: int) -> bool:
if user_id not in self._failed_attempts: if user_id not in self._failed_attempts:
return False return False
+2 -1
View File
@@ -2295,7 +2295,7 @@ class BotInterface:
# Resume < 5 мин: продолжает с N+1 без задержки # Resume < 5 мин: продолжает с N+1 без задержки
# Resume > 5 мин: серия этой ссылки отбрасывается, ждём следующий цикл # Resume > 5 мин: серия этой ссылки отбрасывается, ждём следующий цикл
remaining_clicks = 0 remaining_clicks = params.remaining_clicks # восстанавливаем из сохранённого состояния
series_paused_at: Optional[float] = None series_paused_at: Optional[float] = None
SERIES_EXPIRY = 5 * 60 SERIES_EXPIRY = 5 * 60
@@ -2372,6 +2372,7 @@ class BotInterface:
# Клик завершён — проверяем паузу # Клик завершён — проверяем паузу
if params.paused or params.stopped: if params.paused or params.stopped:
remaining_clicks = series_size - click_num - 1 remaining_clicks = series_size - click_num - 1
params.remaining_clicks = remaining_clicks
if remaining_clicks > 0 and not params.stopped: if remaining_clicks > 0 and not params.stopped:
series_paused_at = datetime.now().timestamp() series_paused_at = datetime.now().timestamp()
logger.info(f"⏸️ Paused after click {click_num+1}/{series_size}, {remaining_clicks} saved") logger.info(f"⏸️ Paused after click {click_num+1}/{series_size}, {remaining_clicks} saved")
+16 -17
View File
@@ -113,25 +113,25 @@ class BotApplication:
async def _restore_tasks(self): async def _restore_tasks(self):
"""Восстанавливает задачи после перезапуска.""" """Восстанавливает задачи после перезапуска."""
saved_tasks = await self.storage.load_tasks() saved_tasks = await self.storage.load_tasks()
if not saved_tasks: if not saved_tasks:
logger.info("No tasks to restore") logger.info("No tasks to restore")
return return
restored = 0 restored = 0
notify_users: dict = {} # user_id -> chat_id # user_id -> (chat_id, was_paused)
notify_users: dict = {}
for task_id, params in saved_tasks.items(): for task_id, params in saved_tasks.items():
# Завершённые и остановленные — загружаем в память но не запускаем # Завершённые и остановленные — только в память, не запускаем
if params.completed or params.stopped: if params.completed or params.stopped:
await self.interface.task_manager.add_task(task_id, params) await self.interface.task_manager.add_task(task_id, params)
continue continue
# Активные задачи ставим на паузу и перезапускаем was_paused = params.paused # сохраняем намерение пользователя
# Если стрим был оффлайн — сохраняем этот статус
params.paused = True
if params.stream_offline: if params.stream_offline:
logger.info(f"Task {task_id}: stream was offline, keeping paused") logger.info(f"Task {task_id}: stream was offline, keeping paused")
await self.interface.task_manager.add_task(task_id, params) await self.interface.task_manager.add_task(task_id, params)
if params.task_type == "twitch_irc" and params.chat_id: if params.task_type == "twitch_irc" and params.chat_id:
@@ -156,20 +156,19 @@ class BotApplication:
restored += 1 restored += 1
if params.user_id and params.chat_id: if params.user_id and params.chat_id:
notify_users[params.user_id] = params.chat_id notify_users[params.user_id] = (params.chat_id, was_paused)
# Уведомляем пользователей о рестарте # Уведомляем пользователей разные сообщения для паузы и активных задач
for uid, chat_id in notify_users.items(): for uid, (chat_id, was_paused) in notify_users.items():
try: try:
await self.bot.send_message( if was_paused:
chat_id, msg = "🔄 Бот был перезапущен\n\n⏸ Ваша задача остаётся на паузе."
"🔄 Бот был перезапущен\n\n" else:
"Ваши задачи приостановлены.\n" msg = "🔄 Бот был перезапущен\n\n▶️ Ваши задачи продолжают работу автоматически."
"Нажмите ▶️ в списке задач для возобновления." await self.bot.send_message(chat_id, msg)
)
except Exception as e: except Exception as e:
logger.warning(f"Failed to notify user {uid}: {e}") logger.warning(f"Failed to notify user {uid}: {e}")
logger.info(f"🔄 Restored {restored} tasks") logger.info(f"🔄 Restored {restored} tasks")
async def start(self): async def start(self):
+2
View File
@@ -54,6 +54,7 @@ class TaskStorage:
"max_ctr": params.max_ctr, "max_ctr": params.max_ctr,
"started_at": params.started_at.isoformat() if params.started_at else None, "started_at": params.started_at.isoformat() if params.started_at else None,
"stream_offline": params.stream_offline, "stream_offline": params.stream_offline,
"remaining_clicks": params.remaining_clicks,
# pending_visits и total_planned_visits — runtime, не сохраняем # pending_visits и total_planned_visits — runtime, не сохраняем
} }
@@ -88,6 +89,7 @@ class TaskStorage:
min_ctr=data.get("min_ctr", 0.8), min_ctr=data.get("min_ctr", 0.8),
max_ctr=data.get("max_ctr", 1.0), max_ctr=data.get("max_ctr", 1.0),
stream_offline=data.get("stream_offline", False), stream_offline=data.get("stream_offline", False),
remaining_clicks=data.get("remaining_clicks", 0),
) )
if data.get("started_at"): if data.get("started_at"):
try: try:
+1
View File
@@ -59,6 +59,7 @@ class TaskParams:
links_found: int = 0 links_found: int = 0
pending_visits: int = 0 # кликов в очереди (runtime, уменьшается) pending_visits: int = 0 # кликов в очереди (runtime, уменьшается)
total_planned_visits: int = 0 # всего запланировано кликов по всем ссылкам (runtime) total_planned_visits: int = 0 # всего запланировано кликов по всем ссылкам (runtime)
remaining_clicks: int = 0 # остаток кликов в текущей серии (сохраняется)
chat_id: Optional[int] = None chat_id: Optional[int] = None
user_id: Optional[int] = None # Telegram user_id, если задача создана пользователем user_id: Optional[int] = None # Telegram user_id, если задача создана пользователем
+336
View File
@@ -0,0 +1,336 @@
"""
Утилита проверки прокси.
- Удаляет дубликаты по ip:port
- Тестирует напрямую (HTTP и SOCKS5)
- Тестирует через локальный SOCKS5→HTTP туннель
- Сохраняет результат обратно в файл
Использование:
python tools/proxy_checker.py
python tools/proxy_checker.py --file proxies.txt --concurrency 10
python tools/proxy_checker.py --no-save # только показать, не писать файл
"""
import asyncio
import sys
import time
import argparse
from pathlib import Path
from dataclasses import dataclass, field
from typing import Optional
sys.path.insert(0, str(Path(__file__).parent.parent))
from managers.proxy_manager import Proxy, ProxyType
TEST_URL = "https://www.google.com/generate_204"
TIMEOUT = 10.0
@dataclass
class CheckResult:
proxy: Proxy
direct_ok: bool = False
direct_ms: Optional[float] = None
tunnel_ok: bool = False
tunnel_ms: Optional[float] = None
error: Optional[str] = None
@property
def ok(self) -> bool:
return self.direct_ok or self.tunnel_ok
def summary(self) -> str:
parts = []
if self.direct_ok:
parts.append(f"direct {self.direct_ms:.0f}ms")
if self.tunnel_ok:
parts.append(f"tunnel {self.tunnel_ms:.0f}ms")
if not parts:
return f"FAIL {self.error or ''}"
return "OK " + " | ".join(parts)
async def _check_direct(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Прямая проверка через requests в thread pool."""
import asyncio
loop = asyncio.get_running_loop()
def _sync():
import requests
t0 = time.perf_counter()
scheme = proxy.proxy_type.value
if proxy.login and proxy.password:
url = f"{scheme}://{proxy.login}:{proxy.password}@{proxy.ip}:{proxy.port}"
else:
url = f"{scheme}://{proxy.ip}:{proxy.port}"
proxies = {"http": url, "https": url}
try:
r = requests.get(TEST_URL, proxies=proxies, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000
return r.status_code in (200, 204), ms
except Exception:
return False, None
return await loop.run_in_executor(None, _sync)
async def _check_tunnel(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Проверка через SOCKS5→HTTP туннель (только для SOCKS5/SOCKS4)."""
if proxy.proxy_type not in (ProxyType.SOCKS5, ProxyType.SOCKS4):
return False, None
from services.socks5_to_http_proxy import Socks5ToHttpProxy
import aiohttp
tunnel = Socks5ToHttpProxy(
socks5_host=proxy.ip,
socks5_port=int(proxy.port),
username=proxy.login,
password=proxy.password,
)
try:
await asyncio.wait_for(tunnel.start(), timeout=5)
t0 = time.perf_counter()
ok = await tunnel.check_connection(test_url=TEST_URL, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000 if ok else None
return ok, ms
except Exception:
return False, None
finally:
try:
await asyncio.wait_for(tunnel.stop(), timeout=3)
except Exception:
pass
async def _check_direct_as_http(proxy: Proxy, timeout: float) -> tuple[bool, Optional[float]]:
"""Принудительная проверка как HTTP прокси (игнорирует настроенный тип)."""
loop = asyncio.get_running_loop()
def _sync():
import requests
t0 = time.perf_counter()
if proxy.login and proxy.password:
url = f"http://{proxy.login}:{proxy.password}@{proxy.ip}:{proxy.port}"
else:
url = f"http://{proxy.ip}:{proxy.port}"
proxies = {"http": url, "https": url}
try:
r = requests.get(TEST_URL, proxies=proxies, timeout=timeout)
ms = (time.perf_counter() - t0) * 1000
return r.status_code in (200, 204), ms
except Exception:
return False, None
return await loop.run_in_executor(None, _sync)
async def check_proxy(proxy: Proxy, timeout: float) -> CheckResult:
result = CheckResult(proxy=proxy)
try:
direct_ok, direct_ms = await asyncio.wait_for(
_check_direct(proxy, timeout), timeout=timeout + 2
)
result.direct_ok = direct_ok
result.direct_ms = direct_ms
if proxy.proxy_type in (ProxyType.SOCKS5, ProxyType.SOCKS4):
tunnel_ok, tunnel_ms = await asyncio.wait_for(
_check_tunnel(proxy, timeout), timeout=timeout + 5
)
result.tunnel_ok = tunnel_ok
result.tunnel_ms = tunnel_ms
except Exception as e:
result.error = str(e)[:60]
return result
async def retry_failed_as_http(results: list[CheckResult], timeout: float, sem: asyncio.Semaphore) -> dict[str, CheckResult]:
"""
Повторно проверяет упавшие прокси принудительно как HTTP.
Возвращает dict[ip:port -> CheckResult] только для тех что теперь прошли.
"""
failed = [r for r in results if not r.ok]
if not failed:
return {}
recovered: dict[str, CheckResult] = {}
async def _retry(r: CheckResult):
async with sem:
try:
ok, ms = await asyncio.wait_for(
_check_direct_as_http(r.proxy, timeout), timeout=timeout + 2
)
except (TimeoutError, asyncio.TimeoutError, asyncio.CancelledError, Exception):
ok, ms = False, None
if ok:
key = f"{r.proxy.ip}:{r.proxy.port}"
r.direct_ok = True
r.direct_ms = ms
r.error = None
r.proxy.proxy_type = ProxyType.HTTP
recovered[key] = r
await asyncio.gather(*[_retry(r) for r in failed], return_exceptions=True)
return recovered
def _load_proxies(file_path: str) -> tuple[list[Proxy], list[str]]:
"""Загружает прокси из файла. Возвращает (список прокси, сырые строки)."""
raw_lines = Path(file_path).read_text(encoding="utf-8").splitlines()
proxies = []
seen = set() # ip:port для дедупликации
for line in raw_lines:
stripped = line.strip()
if not stripped or stripped.startswith("#") or stripped.startswith("!"):
continue
p = Proxy.from_line(stripped)
if not p:
continue
key = f"{p.ip}:{p.port}"
if key in seen:
continue
seen.add(key)
proxies.append(p)
return proxies, raw_lines
def _save_results(file_path: str, results: list[CheckResult], raw_lines: list[str]):
"""Перезаписывает файл: рабочие без '!', нерабочие с '!'.
Использует актуальный тип прокси из results (может быть изменён retry-фазой).
"""
# ip:port → актуальный объект Proxy с возможно обновлённым типом
proxy_map: dict[str, Proxy] = {
f"{r.proxy.ip}:{r.proxy.port}": r.proxy for r in results
}
working = {f"{r.proxy.ip}:{r.proxy.port}" for r in results if r.ok}
new_lines = []
seen_keys: set[str] = set() # для удаления дубликатов
for line in raw_lines:
stripped = line.strip()
if not stripped or stripped.startswith("#"):
new_lines.append(line)
continue
clean = stripped.lstrip("!")
p = Proxy.from_line(clean)
if not p:
new_lines.append(line)
continue
key = f"{p.ip}:{p.port}"
if key in seen_keys:
continue # дубликат — пропускаем
seen_keys.add(key)
actual = proxy_map.get(key, p)
scheme = actual.proxy_type.value
if actual.login and actual.password:
formatted = f"{scheme}://{actual.ip}:{actual.port}:{actual.login}:{actual.password}"
else:
formatted = f"{scheme}://{actual.ip}:{actual.port}"
if key in working:
new_lines.append(formatted)
else:
new_lines.append(f"!{formatted}")
Path(file_path).write_text("\n".join(new_lines) + "\n", encoding="utf-8")
async def run_checker(file_path: str, concurrency: int, timeout: float, save: bool, retry_http: bool):
proxies, raw_lines = _load_proxies(file_path)
total_raw = sum(
1 for l in raw_lines
if l.strip() and not l.strip().startswith("#") and not l.strip().startswith("!")
)
duplicates = total_raw - len(proxies)
print(f"\n{'='*60}")
print(f" Файл: {file_path}")
print(f" Загружено: {total_raw} | Дубликатов: {duplicates} | К проверке: {len(proxies)}")
print(f" Параллельность: {concurrency} | Таймаут: {timeout}s")
if retry_http:
print(f" Повтор упавших через HTTP: ВКЛ")
print(f"{'='*60}\n")
sem = asyncio.Semaphore(concurrency)
results: list[CheckResult] = []
done = 0
async def _run(proxy: Proxy):
nonlocal done
async with sem:
r = await check_proxy(proxy, timeout)
results.append(r)
done += 1
status = "" if r.ok else ""
print(f" [{done:>3}/{len(proxies)}] {status} {proxy.id:<35} {r.summary()}")
await asyncio.gather(*[_run(p) for p in proxies])
# Фаза 2: повторная проверка упавших как HTTP
recovered_count = 0
if retry_http:
failed_count = sum(1 for r in results if not r.ok)
if failed_count:
print(f"\n ── Повтор {failed_count} упавших как HTTP прямую ──\n")
recovered = await retry_failed_as_http(results, timeout, sem)
recovered_count = len(recovered)
for key, r in recovered.items():
print(f" [retry] ✅ {r.proxy.ip}:{r.proxy.port} → HTTP {r.direct_ms:.0f}ms")
ok_count = sum(1 for r in results if r.ok)
fail_count = len(results) - ok_count
direct_ok = sum(1 for r in results if r.direct_ok)
tunnel_ok = sum(1 for r in results if r.tunnel_ok)
print(f"\n{'='*60}")
print(f" ИТОГО")
print(f"{'='*60}")
print(f" Рабочих : {ok_count} / {len(proxies)}")
print(f" Нерабочих : {fail_count}")
print(f" Direct OK : {direct_ok}")
print(f" Tunnel OK : {tunnel_ok}")
if recovered_count:
print(f" Спасено HTTP: {recovered_count}")
if duplicates:
print(f" Дубликатов : {duplicates} (удалены)")
if save:
_save_results(file_path, results, raw_lines)
print(f"\n Файл обновлён: {file_path}")
print(f" Нерабочие помечены '!'")
print(f"{'='*60}\n")
def main():
parser = argparse.ArgumentParser(description="Проверка прокси")
parser.add_argument("--file", default="proxies.txt", help="Файл с прокси")
parser.add_argument("--concurrency", type=int, default=5, help="Параллельность (по умолчанию: 5)")
parser.add_argument("--timeout", type=float, default=10.0, help="Таймаут секунд (по умолчанию: 10)")
parser.add_argument("--no-save", action="store_true", help="Не перезаписывать файл")
parser.add_argument("--retry-http", action="store_true", help="Повторить упавшие как HTTP напрямую")
args = parser.parse_args()
asyncio.run(run_checker(
file_path=args.file,
concurrency=args.concurrency,
timeout=args.timeout,
save=not args.no_save,
retry_http=args.retry_http,
))
if __name__ == "__main__":
main()