"""Главный цикл: обход разделов через адаптер. Единая точка запуска для любого движка форума. Не содержит кода парсинга — использует адаптер из forums/. """ import sys import logging import shutil import json from concurrent.futures import ThreadPoolExecutor, as_completed from pathlib import Path from . import fetch as _fetch from .throttle import Throttle, BanDetector from .output import TopicWriter from .state import load, save, merge_workers log = logging.getLogger(__name__) # ─── Глобальный throttle (виден fetch.py) ─────────────────────────────── _throttle: Throttle | None = None # ═══════════════════════════════════════════════════════════════════════════ # Функция одного воркера (запускается в ThreadPoolExecutor) # ═══════════════════════════════════════════════════════════════════════════ def _worker(section_url: str, section_slug: str, section_name: str, pages: range, worker_id: int, adapter) -> tuple: """Обойти свой диапазон страниц и сохранить темы в worker_N/. Args: section_url: URL раздела (первая страница). section_slug: Короткое имя для папки (напр. "diagnostika"). section_name: Человеческое название раздела. pages: Диапазон страниц для этого воркера. worker_id: Номер воркера (0, 1, 2...). adapter: Экземпляр адаптера (forums/*/adapter.py). Returns: (worker_id, topic_count) — для汇总 в главном потоке. """ session = _fetch.make_session() # ── Состояние воркера ─────────────────────────────────────────────── worker_state = { "topics_done": {}, "pages_done": [], "file_index": 0, "topic_count": 0, } consecutive_404 = 0 writer = None # ── Обход страниц раздела ─────────────────────────────────────────── for page_num in pages: # Формируем URL страницы через адаптер url = adapter.page_url(section_url, page_num) log.info("[W%d] 📄 стр.%d: %s", worker_id, page_num, url[:100]) # Загружаем soup = _fetch.get(url, session) if soup is None: # 404 подряд — вероятно, страницы кончились consecutive_404 += 1 if consecutive_404 >= 3: log.info("[W%d] ⏹️ 3 пустых страницы подряд — завершаем", worker_id) break continue consecutive_404 = 0 # Парсим список тем через адаптер topics = adapter.parse_topics(soup) log.info("[W%d] Тем на странице: %d", worker_id, len(topics)) # ── Обход тем на странице ──────────────────────────────────────── for idx, topic in enumerate(topics, 1): topic_id = str(topic["id"]) # Пропускаем уже обработанные if topic_id in worker_state["topics_done"]: continue log.info( "[W%d] [%d/%d] #%s: %s", worker_id, idx, len(topics), topic_id, topic["title"][:80], ) # Загружаем страницу темы (первую) topic_url = adapter.topic_url(section_url, topic) topic["url"] = topic_url topic_soup = _fetch.get(topic_url, session) if topic_soup is None: # Не загрузилась — помечаем как сделанную и идём дальше worker_state["topics_done"][topic_id] = topic["title"] continue # Определяем количество страниц в теме topic_pages = adapter.count_pages(topic_soup) log.info( "[W%d] Страниц в теме: %d", worker_id, topic_pages, ) # Парсим посты с первой страницы all_posts = adapter.parse_posts(topic_soup) # Догружаем остальные страницы темы for tp in range(2, topic_pages + 1): tp_url = adapter.topic_page_url(topic_url, tp) tp_soup = _fetch.get(tp_url, session) if tp_soup: all_posts.extend(adapter.parse_posts(tp_soup)) # Парсим теги (если адаптер поддерживает) tags = adapter.parse_tags(topic_soup) if hasattr(adapter, "parse_tags") else [] log.info( "[W%d] %d постов, теги: %s", worker_id, len(all_posts), tags, ) # ── Сохраняем ──────────────────────────────────────────────── # Лениво инициализируем writer при первой теме if writer is None: worker_dir = Path(adapter.data_dir) / section_slug / f"worker_{worker_id}" worker_dir.mkdir(parents=True, exist_ok=True) writer = TopicWriter(worker_dir, adapter.topics_per_file) writer.write(section_name, section_slug, topic, all_posts, tags) worker_state["topics_done"][topic_id] = topic["title"] worker_state["topic_count"] += 1 # Помечаем страницу как пройденную worker_state["pages_done"].append(page_num) # ── Сохраняем состояние воркера ───────────────────────────────────── if writer is not None: writer.close() worker_dir = Path(adapter.data_dir) / section_slug / f"worker_{worker_id}" worker_dir.mkdir(parents=True, exist_ok=True) state_file = worker_dir / "state.json" state_file.write_text( json.dumps(worker_state, ensure_ascii=False, indent=2), encoding="utf-8", ) session.close() return worker_id, worker_state["topic_count"] # ═══════════════════════════════════════════════════════════════════════════ # Функция мержа воркеров # ═══════════════════════════════════════════════════════════════════════════ def _merge_workers(adapter, section_slug: str): """Собрать part_*.jsonl из worker_N/ в корень раздела. Переименовывает файлы с единой нумерацией, удаляет пустые worker-папки, мержит state.json. Args: adapter: Экземпляр адаптера. section_slug: Короткое имя раздела. """ root = Path(adapter.data_dir) / section_slug # ── Собираем все part_*.jsonl из worker_N/ ────────────────────────── all_parts = [] for worker_dir in sorted(root.glob("worker_*")): if worker_dir.is_dir(): for part_file in sorted(worker_dir.glob("part_*.jsonl")): all_parts.append(part_file) if not all_parts: log.warning("⚠️ Нет файлов для мержа в %s", root) return log.info( "🔗 Мерж %d файлов из %d воркеров → %s", len(all_parts), len(list(root.glob("worker_*"))), root, ) # ── Переименовываем с единой нумерацией ───────────────────────────── for new_index, src in enumerate(all_parts): dst = root / f"part_{new_index:04d}.jsonl" shutil.move(str(src), str(dst)) # ── Удаляем пустые папки воркеров ─────────────────────────────────── for worker_dir in root.glob("worker_*"): if worker_dir.is_dir(): try: worker_dir.rmdir() # удаляет только пустую папку except OSError: pass # если не пустая — оставляем # ── Мержим state.json ─────────────────────────────────────────────── merge_workers(adapter.data_dir, section_slug) log.info("✅ Мерж завершён: %d файлов → %s", len(all_parts), root) # ═══════════════════════════════════════════════════════════════════════════ # Публичный запуск # ═══════════════════════════════════════════════════════════════════════════ def run(adapter, section_url: str, section_slug: str | None = None, workers: int = 4, no_delay: bool = False, ban_test: bool = True): """Главный цикл: многопоточный обход раздела через адаптер. Порядок работы: 1. Проверка бана (BanDetector). 2. Загрузка первой страницы для определения названия и пагинации. 3. Разбиение страниц на диапазоны по числу воркеров. 4. Запуск ThreadPoolExecutor. 5. Мерж результатов. Args: adapter: Экземпляр адаптера (forums/*/adapter.py). section_url: URL раздела (первая страница). section_slug: Короткое имя (если None — берётся из URL). workers: Количество потоков. no_delay: True — отключить координацию задержек. ban_test: True — проверить бан перед стартом. """ # ── Определяем slug из URL ────────────────────────────────────────── if section_slug is None: section_slug = section_url.rstrip("/").split("/")[-1] # ── Настройка throttle ────────────────────────────────────────────── throttle_enabled = not no_delay # Если workers=1, координация не нужна — один поток не конфликтует if workers <= 1: throttle_enabled = False # Проверка бана if ban_test and throttle_enabled and workers > 1: if BanDetector.test(section_url, adapter.ban_test_count): throttle_enabled = True else: throttle_enabled = False # Устанавливаем глобальный throttle global _throttle _throttle = Throttle(adapter.delay, enabled=throttle_enabled) _fetch.throttle = _throttle log.info( "⚙️ Потоков: %d, координация: %s", workers, "вкл" if throttle_enabled else "выкл", ) # ── Определяем название раздела и число страниц ───────────────────── session = _fetch.make_session() soup = _fetch.get(section_url, session) if soup is None: log.error("❌ Сайт недоступен: %s", section_url) sys.exit(1) # Название раздела (из h1 или из URL) section_name = adapter.parse_section_name(soup) or section_slug # Количество страниц в разделе total_pages = adapter.count_pages(soup) log.info("📋 %s | страниц: %d", section_name, total_pages) session.close() # ── Разбиваем страницы на диапазоны для воркеров ──────────────────── if workers > total_pages: workers = max(1, 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("📦 Диапазоны: %s", [f"{r.start}-{r[-1]}" for r in ranges]) # ── Запуск воркеров ───────────────────────────────────────────────── total_topics = 0 with ThreadPoolExecutor(max_workers=workers) as pool: futures = [ pool.submit(_worker, section_url, section_slug, section_name, rng, wid, adapter) for wid, rng in enumerate(ranges) ] for future in as_completed(futures): wid, cnt = future.result() total_topics += cnt log.info("[W%d] ✅ завершён: %d тем", wid, cnt) # ── Мерж ──────────────────────────────────────────────────────────── _merge_workers(adapter, section_slug) # ── Итог ──────────────────────────────────────────────────────────── data_dir = Path(adapter.data_dir) part_files = sorted((data_dir / section_slug).glob("part_*.jsonl")) total_size_mb = sum(f.stat().st_size for f in part_files) / (1024 * 1024) log.info("=" * 60) log.info( "✅ '%s' — %d тем в %d частях (%.1f MB)", section_name, total_topics, len(part_files), total_size_mb, ) log.info("=" * 60)