""" Сохранение и загрузка задач между перезапусками. """ import asyncio import json import logging from pathlib import Path from datetime import datetime from typing import Dict from managers.task_manager import TaskParams logger = logging.getLogger(__name__) class TaskStorage: """Хранилище задач в JSON файле.""" def __init__(self, file_path: str = "data/tasks.json"): self.file_path = Path(file_path) self.file_path.parent.mkdir(parents=True, exist_ok=True) self._lock = asyncio.Lock() def _params_to_dict(self, params: TaskParams) -> dict: return { "url": params.url, "task_type": params.task_type, "min_delay": params.min_delay, "max_delay": params.max_delay, "min_reading": params.min_reading, "max_reading": params.max_reading, "max_visits": params.max_visits, "current_visit": params.current_visit, "channel": params.channel, "target_username": params.target_username, "visits_per_link": params.visits_per_link, "monitor_minutes": params.monitor_minutes, "allowed_domains": params.allowed_domains, "paused": params.paused, "completed": params.completed, "stopped": params.stopped, "total_visits": params.total_visits, "successful_visits": params.successful_visits, "links_found": params.links_found, "chat_id": params.chat_id, "started_at": params.started_at.isoformat() if params.started_at else None, } def _dict_to_params(self, data: dict) -> TaskParams: params = TaskParams( url=data.get("url", ""), task_type=data.get("task_type", "visit"), min_delay=data.get("min_delay", 10), max_delay=data.get("max_delay", 30), min_reading=data.get("min_reading", 5), max_reading=data.get("max_reading", 15), max_visits=data.get("max_visits"), current_visit=data.get("current_visit", 0), channel=data.get("channel", ""), target_username=data.get("target_username", "*"), visits_per_link=data.get("visits_per_link", 5), monitor_minutes=data.get("monitor_minutes", 10), allowed_domains=data.get("allowed_domains"), paused=data.get("paused", False), completed=data.get("completed", False), stopped=data.get("stopped", False), total_visits=data.get("total_visits", 0), successful_visits=data.get("successful_visits", 0), links_found=data.get("links_found", 0), chat_id=data.get("chat_id"), ) if data.get("started_at"): try: params.started_at = datetime.fromisoformat(data["started_at"]) except (ValueError, TypeError) as e: logger.warning(f"Failed to parse started_at: {e}") return params def _save_sync(self, data: dict) -> None: with open(self.file_path, "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False, indent=2) def _load_sync(self) -> dict: if not self.file_path.exists(): return {} with open(self.file_path, "r", encoding="utf-8") as f: return json.load(f) async def save_tasks(self, tasks: Dict[str, TaskParams]) -> None: async with self._lock: try: data = {tid: self._params_to_dict(p) for tid, p in tasks.items()} await asyncio.to_thread(self._save_sync, data) logger.info(f"💾 Saved {len(data)} tasks to {self.file_path}") except Exception as e: logger.error(f"Save error: {e}") async def load_tasks(self) -> Dict[str, TaskParams]: try: data = await asyncio.to_thread(self._load_sync) if not data: logger.info("No saved tasks found") return {} tasks = {tid: self._dict_to_params(d) for tid, d in data.items()} logger.info(f"📂 Loaded {len(tasks)} tasks from {self.file_path}") return tasks except Exception as e: logger.error(f"Load error: {e}") return {} async def save_task(self, task_id: str, params: TaskParams) -> None: tasks = await self.load_tasks() tasks[task_id] = params await self.save_tasks(tasks) async def delete_task(self, task_id: str) -> None: tasks = await self.load_tasks() if task_id in tasks: del tasks[task_id] await self.save_tasks(tasks) class ChatStorage: """Хранилище привязки стримеров к чатам.""" def __init__(self, file_path: str = "data/chats.json"): self.file_path = Path(file_path) self.file_path.parent.mkdir(parents=True, exist_ok=True) self._lock = asyncio.Lock() def _load_sync(self) -> dict: if not self.file_path.exists(): return {} try: with open(self.file_path, "r", encoding="utf-8") as f: return json.load(f) except (json.JSONDecodeError, IOError) as e: logger.warning(f"Failed to load chats: {e}") return {} def _save_sync(self, data: dict) -> None: with open(self.file_path, "w", encoding="utf-8") as f: json.dump(data, f, ensure_ascii=False, indent=2) async def get_chat_streamers(self, chat_id: int) -> list: async with self._lock: data = await asyncio.to_thread(self._load_sync) return data.get(str(chat_id), []) async def add_streamer(self, chat_id: int, channel: str) -> None: async with self._lock: data = await asyncio.to_thread(self._load_sync) key = str(chat_id) if key not in data: data[key] = [] if channel not in data[key]: data[key].append(channel) await asyncio.to_thread(self._save_sync, data) async def remove_streamer(self, chat_id: int, channel: str) -> None: async with self._lock: data = await asyncio.to_thread(self._load_sync) key = str(chat_id) if key in data and channel in data[key]: data[key].remove(channel) await asyncio.to_thread(self._save_sync, data)