diff --git a/scraper_jobs.py b/scraper_jobs.py index 2b50bde..216d981 100644 --- a/scraper_jobs.py +++ b/scraper_jobs.py @@ -1,6 +1,6 @@ import asyncio import logging -from typing import Any, Dict, List +from typing import Any, Dict, List, Optional from app_state import StateStore @@ -25,12 +25,21 @@ class ScraperJobService: asyncio.run(self._run_async(job_type, payload)) async def _run_async(self, job_type: str, payload: Dict[str, Any]) -> None: + # Extract account_id from payload, default to None (legacy) + account_id: Optional[str] = payload.get("account_id") ScraperClass = self._import_scraper_class() - scraper = ScraperClass() + scraper = ScraperClass(account_id=account_id) + + if account_id: + from app_state import load_account + acc_state = load_account(self.state_store.path.parent, account_id) + # Load per-account state into scraper state + scraper.state = acc_state + initialized = await scraper.initialize_client(interactive=False) if not initialized: raise RuntimeError( - "Telegram client is not ready. Check credentials, session, and write access to /app/session." + "Telegram client is not ready. Check credentials, session, and write access." ) try: if job_type == "scrape_channel": diff --git a/telegram_scraper_with_forwarding.py b/telegram_scraper_with_forwarding.py index b6bfe51..c60f676 100644 --- a/telegram_scraper_with_forwarding.py +++ b/telegram_scraper_with_forwarding.py @@ -21,7 +21,14 @@ from telethon.tl.types import ( ) from telethon.errors import FloodWaitError, SessionPasswordNeededError import qrcode -from app_state import StateStore +from app_state import ( + StateStore, + account_data_dir, + account_session_path, + get_account_store, + load_account, + save_account, +) warnings.filterwarnings( "ignore", message="Using async sessions support is an experimental feature" @@ -73,13 +80,21 @@ class ForwardingRule: class OptimizedTelegramScraper: - def __init__(self): - self.DATA_DIR = Path("data") + def __init__(self, account_id: Optional[str] = None): + self.account_id = account_id self.SESSION_DIR = Path("session") - self.DATA_DIR.mkdir(exist_ok=True) self.SESSION_DIR.mkdir(exist_ok=True) + + if account_id: + self.DATA_DIR = Path("data") / "accounts" / account_id + self.state_store = get_account_store(Path("data"), account_id) + else: + self.DATA_DIR = Path("data") + self.state_store = StateStore(self.DATA_DIR / "state.json") + + self.DATA_DIR.mkdir(parents=True, exist_ok=True) self.STATE_FILE = str(self.DATA_DIR / "state.json") - self.state_store = StateStore(self.DATA_DIR / "state.json") + self.state = self.load_state() self.client = None self.continuous_scraping_active = False @@ -1480,8 +1495,9 @@ class OptimizedTelegramScraper: print("Invalid API ID. Must be a number.") return False + session = account_session_path(self.SESSION_DIR, self.account_id or "session") self.client = TelegramClient( - str(self.SESSION_DIR / "session"), + session, self.state["api_id"], self.state["api_hash"], )