"""Главный цикл: обход разделов phpBB-форума.""" import sys, logging, shutil, json, re from concurrent.futures import ThreadPoolExecutor, as_completed from pathlib import Path from urllib.parse import urljoin from .config import ( OUTPUT_DIR, WORKERS, BAN_TEST, DELAY, SECTIONS, FORUM_BASE, TOPICS_PER_PAGE, POSTS_PER_PAGE, ) from .fetch import get, make_session from .paginate import count from .parse import parse_topics, parse_posts from .state import load, save, merge_states from .output import save_topic from .throttle import Throttle, BanDetector from . import fetch as _fetch log = logging.getLogger(__name__) _throttle = None def _url(path: str) -> str: """Полный URL относительно FORUM_BASE.""" if path.startswith("http"): # Убираем sid из URL для чистоты path = re.sub(r"[?&]sid=[^&]+", "", path) # Заменяем http на https path = path.replace("http://", "https://") return path path = re.sub(r"[?&]sid=[^&]+", "", path) return urljoin(FORUM_BASE, path) def _worker(forum_id: int, section_name: str, pages: range, worker_id: int): """Один поток: обходит свой диапазон страниц раздела.""" session = make_session() st = {"topics_done": {}, "pages_done": [], "file_index": 0, "topic_count": 0} slug = f"f{forum_id}" consecutive_404 = 0 worker_dir = OUTPUT_DIR / slug / f"worker_{worker_id}" worker_dir.mkdir(parents=True, exist_ok=True) base_url = f"{FORUM_BASE}viewforum.php?f={forum_id}" for pg in pages: url = base_url if pg == 1 else f"{base_url}&start={(pg-1)*TOPICS_PER_PAGE}" log.info(f"[W{worker_id}] 📄 стр.{pg}") page_soup = get(url, session) if not page_soup: consecutive_404 += 1 if consecutive_404 >= 3: log.info(f"[W{worker_id}] ⏹️ 3 пустых — завершаем") break continue consecutive_404 = 0 topics = parse_topics(page_soup) log.info(f"[W{worker_id}] Тем: {len(topics)}") for i, tp in enumerate(topics, 1): tid = tp["id"] if str(tid) in st["topics_done"]: continue # Собираем полный URL темы topic_url = _url(tp["url"] or f"./viewtopic.php?f={forum_id}&t={tid}") tp["url"] = topic_url log.info(f"[W{worker_id}] [{i}/{len(topics)}] #{tid}: {tp['title'][:80]}") topic_soup = get(topic_url, session) if not topic_soup: st["topics_done"][str(tid)] = tp["title"] continue topic_pages = count(topic_soup) posts = [] for tpp in range(1, topic_pages + 1): if tpp == 1: posts += parse_posts(topic_soup) else: tp_url = f"{topic_url}&start={(tpp-1)*POSTS_PER_PAGE}" tp2 = get(tp_url, session) if tp2: posts += parse_posts(tp2) log.info(f"[W{worker_id}] {len(posts)} постов ({topic_pages} стр.)") save_topic(slug, section_name, tp, posts, st, worker_dir=worker_dir) st["topics_done"][str(tid)] = tp["title"] st["pages_done"].append(pg) _state_file = worker_dir / "state.json" _state_file.write_text(json.dumps(st, ensure_ascii=False, indent=2), encoding="utf-8") session.close() return worker_id, st["topic_count"] def _merge_workers(slug: str): """Собрать part_*.jsonl из всех worker_N/ в корень.""" root = OUTPUT_DIR / slug all_parts = [] for wdir in sorted(root.glob("worker_*")): if wdir.is_dir(): for f in sorted(wdir.glob("part_*.jsonl")): all_parts.append(f) log.info(f"🔗 Мерж {len(all_parts)} файлов из {len(list(root.glob('worker_*')))} воркеров") for idx, src in enumerate(all_parts): dst = root / f"part_{idx:04d}.jsonl" shutil.move(str(src), str(dst)) for wdir in root.glob("worker_*"): if wdir.is_dir(): try: wdir.rmdir() except OSError: pass merge_states(slug) log.info(f"✅ Мерж: {len(all_parts)} файлов → {root}") def run(forum_id: int, workers: int = WORKERS, no_delay: bool = False, ban_test: bool = True): """Главный цикл. Args: forum_id: ID раздела на форуме (напр. 22) workers: количество потоков no_delay: без координации задержек ban_test: проверить бан перед стартом """ slug = f"f{forum_id}" section_name = SECTIONS.get(forum_id, f"Форум #{forum_id}") throttle_enabled = not no_delay if ban_test and throttle_enabled and workers > 1: if not BanDetector.test(f"{FORUM_BASE}viewforum.php?f={forum_id}", BAN_TEST): throttle_enabled = False global _throttle _throttle = Throttle(DELAY, enabled=throttle_enabled) _fetch.throttle = _throttle log.info(f"⚙️ Потоков: {workers}, координация: {'вкл' if throttle_enabled else 'выкл'}") # ── Страница 1: определяем количество страниц ───────────────────── first_url = f"{FORUM_BASE}viewforum.php?f={forum_id}" soup = get(first_url) if not soup: log.error(f"❌ Раздел f={forum_id} недоступен") sys.exit(1) total_pages = count(soup) log.info(f"📋 '{section_name}' (f={forum_id}) | страниц: {total_pages}") if total_pages <= 1: log.info("ℹ️ Всего одна страница — проверяем, есть ли темы...") topics = parse_topics(soup) if not topics: log.warning(f"⚠️ В разделе f={forum_id} нет тем (возможно, это категория с подразделами)") sys.exit(0) # ── Диапазоны страниц ───────────────────────────────────────────── if workers > total_pages: workers = total_pages base_size = total_pages // workers remainder = total_pages % workers ranges = [] start = 1 for w in range(workers): size = base_size + (1 if w < remainder else 0) ranges.append(range(start, start + size)) start += size log.info(f"📦 Диапазоны: {[f'{r.start}-{r[-1]}' for r in ranges]}") # ── Запуск потоков ──────────────────────────────────────────────── total_topics = 0 with ThreadPoolExecutor(max_workers=workers) as pool: futures = [ pool.submit(_worker, forum_id, section_name, rng, i) for i, rng in enumerate(ranges) ] for fut in as_completed(futures): wid, cnt = fut.result() total_topics += cnt log.info(f"[W{wid}] ✅ завершён: {cnt} тем") # ── Мерж ────────────────────────────────────────────────────────── _merge_workers(slug) # ── Итог ────────────────────────────────────────────────────────── log.info("=" * 60) log.info(f"✅ '{section_name}' (f={forum_id}) — {total_topics} тем") for f in sorted((OUTPUT_DIR / slug).glob("part_*.jsonl")): log.info(f" 📄 {f.name} ({f.stat().st_size/1024/1024:.1f} MB)") log.info("=" * 60)