208 lines
8.0 KiB
Python
208 lines
8.0 KiB
Python
"""Главный цикл: обход разделов 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)
|