# VESTI — RSS/Atom-краулер (feedparser + httpx). # Читает источники из sources.yaml (crawler: rss), условные GET (etag/modified), # нормализует URL (UTM/якоря), дотягивает полный текст (trafilatura) для summary-only, # пишет посты в posts (sha256-дедуп) и состояние ленты в rss_state. # # Запуск: # python -m crawler.rss_crawler --all # все включённые rss-источники # python -m crawler.rss_crawler --direction linux # python -m crawler.rss_crawler --source lwn # python -m crawler.rss_crawler --source lwn --dry-run # печать без записи в БД import argparse import hashlib import html as html_mod import logging import re import sys import time from pathlib import Path from urllib.parse import parse_qsl, urlencode, urlparse, urlunparse sys.path.insert(0, str(Path(__file__).resolve().parent.parent)) from config import RSS_DEFAULT_TIMEOUT, RSS_USER_AGENT, ensure_dirs from db.db import db from sources.sources import ( get_enabled_sources, get_source, load_sources_yaml, sync_sources_to_db, start_run, finish_run, ) logging.basicConfig(level=logging.INFO, format="%(asctime)s %(levelname)s %(name)s %(message)s") log = logging.getLogger("vesti.rss") MAX_ITEMS = 25 # максимум записей ленты за запуск (остальные — следующими запусками) RATE_LIMIT_SEC = 1.5 # пауза между источниками (вежливый опрос) MIN_SUMMARY_LEN = 220 # если выжимка короче — пробуем дотянуть полный текст ERROR_DEAD_THRESHOLD = 5 # N ошибок подряд → status='dead' (фид-здоровье) # Параметры запроса, которые режут из URL (трекинг/мусор) UTM_PARAMS = { "utm_source", "utm_medium", "utm_campaign", "utm_term", "utm_content", "fbclid", "gclid", "ref", "ref_src", "mc_cid", "mc_eid", } def normalize_url(url): """Чистит URL от UTM-параметров, fbclid/gclid/ref и якоря -> canonical_url. Возвращает None для мусорных ссылок (без схемы/хоста). """ if not url: return None url = url.strip() if not url.startswith(("http://", "https://")): return None u = urlparse(url) if not u.netloc: return None q = [(k, v) for k, v in parse_qsl(u.query, keep_blank_values=True) if k.lower() not in UTM_PARAMS] return urlunparse(u._replace(query=urlencode(q), fragment="")) def sha256_text(text: str) -> str: return hashlib.sha256((text or "").encode("utf-8", "ignore")).hexdigest() # Блочные теги, вокруг которых сохраняем переносы (абзацы/списки/заголовки/цитаты) _BLOCK_TAGS = [ "p", "div", "li", "h1", "h2", "h3", "h4", "h5", "h6", "blockquote", "pre", "tr", "section", "article", ] # Нормализация текста: сохраняем структуру абзацев (\n\n), убираем мусор. # НЕ схлопываем всё в одну строку — иначе новость в предпросмотре — сплошной текст. def normalize_text(text): """Чистит текст: пробелы по краям строк, 3+ переноса → 2, без изменения слов.""" if not text: return "" text = text.replace("\r\n", "\n").replace("\r", "\n") text = re.sub(r"[ \t]+\n", "\n", text) # пробелы перед концом строки text = re.sub(r"\n[ \t]+", "\n", text) # пробелы после начала строки text = re.sub(r"\n{3,}", "\n\n", text) # не больше одной пустой строки # пробелы по краям каждой строки text = "\n".join(l.strip() for l in text.split("\n")) text = re.sub(r"\n{3,}", "\n\n", text) return text.strip() def html_to_text(html_str): """Раскрывает HTML (из ленты) в читаемый text с сохранением абзацев/списков. HTML-теги разметки (b/i/a/...) вырезаются, их содержимое остаётся текстом; блочные теги (p/li/h1-6/blockquote/pre) дают переносы строк, чтобы текст не превращался в сплошную кашу. """ from selectolax.parser import HTMLParser if not html_str: return "" try: tree = HTMLParser(html_str) tree.strip_tags(["script", "style", "noscript"]) if tree.body is None: return "" for node in tree.css(", ".join(_BLOCK_TAGS)): node.insert_before("\n\n") text = tree.body.text() return normalize_text(text) except Exception as e: log.debug("html_to_text failed: %s", e) return re.sub(r"<[^>]+>", " ", html_mod.unescape(html_str)).strip() def extract_full_text(url, timeout=15): """Полный текст страницы через trafilatura. Возвращает str или None (не падает).""" import httpx try: r = httpx.get(url, headers={"User-Agent": RSS_USER_AGENT}, timeout=timeout, follow_redirects=True) if r.status_code != 200: return None import trafilatura txt = trafilatura.extract( r.text, include_comments=False, include_tables=False, favor_precision=True, # include_formatting=True сохраняет абзацы (\\n\\n) — иначе текст # схлопывается в сплошную строку и новость нечитаема include_formatting=True, ) return (txt or "").strip() or None except Exception as e: log.debug("full text %s failed: %s", url, e) return None def entry_published_at(entry): """published_parsed/updated_parsed -> ISO-строка; None если нет даты.""" for key in ("published_parsed", "updated_parsed"): st = getattr(entry, key, None) if st: try: import datetime return datetime.datetime(*st[:6], tzinfo=datetime.timezone.utc).isoformat() except Exception: return None return None def parse_entries(feed, source, dry_run=False): """Превращает записи ленты feedparser в список dict постов. returns (posts, last_build_date) """ posts = [] for entry in feed.entries[:MAX_ITEMS]: link = getattr(entry, "link", None) or "" canonical = normalize_url(link) if not canonical: continue title = getattr(entry, "title", "") or "" # полный текст из ленты (content[].value), иначе summary content_parts = [c.get("value", "") for c in (getattr(entry, "content", None) or [])] raw_html = content_parts[0] if content_parts else "" summary_html = getattr(entry, "summary", "") or "" text = html_to_text(raw_html or summary_html) # summary-only лента: пробуем дотянуть полный текст if (not raw_html or len(text) < MIN_SUMMARY_LEN) and canonical: full = extract_full_text(canonical) if full: text = full if not text: text = title # последний рубеж: заголовок text = normalize_text(text) posts.append({ "url": canonical, "title": title, "text": text, "published_at": entry_published_at(entry), "author": getattr(entry, "author", None) or None, }) last_build = None fb = getattr(feed, "last_build_date_parsed", None) or getattr(feed, "published_parsed", None) if fb: try: import datetime last_build = datetime.datetime(*fb[:6], tzinfo=datetime.timezone.utc).isoformat() except Exception: last_build = None return posts, last_build def store_rss_posts(conn, source_id, posts): """Вставка RSS-постов в posts с дедупом по sha256(text) и url. Возвращает количество НОВЫХ. Повторная вставка той же статьи (даже с другим URL в ленте) не создаёт второй строки — sha256 от текста. """ new = 0 for p in posts: sha = sha256_text(p["text"]) dup = conn.execute( "SELECT id FROM posts WHERE sha256=? OR url=?", (sha, p["url"]), ).fetchone() if dup: continue conn.execute( """INSERT INTO posts (sha256, source_id, url, text, content_type, published_at, status, is_own, is_own_canonical) VALUES (?,?,?,?, 'text', ?, 'new', 0, 0)""", (sha, source_id, p["url"], p["text"], p["published_at"]), ) new += 1 return new def run_source(source, dry_run=False): """Краулинг одного RSS-источника: условный GET, парсинг, запись, состояние.""" slug = source["slug"] feed_url = source.get("feed_url") if not feed_url: log.error("SKIP %s: нет feed_url (crawler=rss, а лента не задана)", slug) return run_id = start_run(source["id"], f"rss:{slug}", "cron") t0 = time.monotonic() try: import httpx # --- состояние ленты --- with db() as conn: st = conn.execute("SELECT * FROM rss_state WHERE slug=?", (slug,)).fetchone() state = dict(st) if st else {} headers = {"User-Agent": RSS_USER_AGENT, "Accept": "application/rss+xml, application/atom+xml, application/xml;q=0.9, */*;q=0.8"} if state.get("etag"): headers["If-None-Match"] = state["etag"] if state.get("modified"): headers["If-Modified-Since"] = state["modified"] resp = httpx.get(feed_url, headers=headers, timeout=RSS_DEFAULT_TIMEOUT, follow_redirects=True) # вежливый ретрай на 429/503 (LWN троттлит частые опросы) for attempt in range(3): if resp.status_code in (429, 503): retry_after = 0 ra = resp.headers.get("retry-after") if ra: try: retry_after = int(ra) except ValueError: retry_after = 10 wait = max(retry_after, 5 * (attempt + 1)) log.warning("%s: HTTP %d, retry %d/3 через %ds", slug, resp.status_code, attempt + 1, wait) time.sleep(wait) resp = httpx.get(feed_url, headers=headers, timeout=RSS_DEFAULT_TIMEOUT, follow_redirects=True) else: break if resp.status_code == 304: # лента не менялась — экономим (0 fetched/0 new, ok) duration = int((time.monotonic() - t0) * 1000) if not dry_run: finish_run(run_id, "ok", 0, 0, None, duration) log.info("OK %s: 304 (не менялась), %dms", slug, duration) return if resp.status_code != 200: raise RuntimeError(f"HTTP {resp.status_code} for {feed_url}") import feedparser feed = feedparser.parse(resp.content, response_headers=resp.headers) if getattr(feed, "bozo", False) and not feed.entries: raise RuntimeError(f"bozo feed (не парсится): {getattr(feed, 'bozo_exception', '')}") posts, last_build = parse_entries(feed, source, dry_run) if dry_run: print(f"=== {slug} ({source['name']}) — dry_run, {len(posts)} записей ===") for p in posts: print(f" - {p['published_at'] or '????'} | {p['title'][:80]}") print(f" {p['url']}") return with db() as conn: new = store_rss_posts(conn, source["id"], posts) new_etag = (resp.headers.get("etag") or "").strip() new_modified = (resp.headers.get("last-modified") or "").strip() conn.execute( """INSERT INTO rss_state (slug, etag, modified, last_build_date, last_error, error_count, status, updated_at) VALUES (?,?,?,?, NULL, 0, 'alive', datetime('now')) ON CONFLICT(slug) DO UPDATE SET etag=excluded.etag, modified=excluded.modified, last_build_date=excluded.last_build_date, last_error=NULL, error_count=0, status='alive', updated_at=datetime('now')""", (slug, new_etag or state.get("etag"), new_modified or state.get("modified"), last_build), ) conn.execute( "UPDATE sources SET last_fetch=datetime('now'), status='alive', last_error=NULL WHERE id=?", (source["id"],), ) duration = int((time.monotonic() - t0) * 1000) finish_run(run_id, "ok", len(posts), new, None, duration) log.info("OK %s: fetched=%d new=%d (%.1fs)", slug, len(posts), new, duration / 1000) except Exception as e: duration = int((time.monotonic() - t0) * 1000) if not dry_run: with db() as conn: st = conn.execute("SELECT error_count FROM rss_state WHERE slug=?", (slug,)).fetchone() err_count = (st[0] if st else 0) + 1 status = "dead" if err_count >= ERROR_DEAD_THRESHOLD else "alive" conn.execute( """INSERT INTO rss_state (slug, last_error, error_count, status, updated_at) VALUES (?,?,?,?, datetime('now')) ON CONFLICT(slug) DO UPDATE SET last_error=excluded.last_error, error_count=excluded.error_count, status=excluded.status, updated_at=datetime('now')""", (slug, str(e)[:500], err_count, status), ) conn.execute( "UPDATE sources SET last_error=?, status=? WHERE id=?", (str(e)[:500], status, source["id"]), ) finish_run(run_id, "error", 0, 0, str(e)[:500], duration) log.error("ERROR %s: %s", slug, e) def main(): ap = argparse.ArgumentParser(description="VESTI RSS/Atom crawler") ap.add_argument("--source", help="slug источника (например lwn)") ap.add_argument("--direction", help="направление (например linux)") ap.add_argument("--all", action="store_true", help="все включённые rss-источники") ap.add_argument("--dry-run", action="store_true", help="печать записей без записи в БД") args = ap.parse_args() ensure_dirs() data = load_sources_yaml() sync_sources_to_db(data) if args.source: src = get_source(args.source) if not src: log.error("source %s not found", args.source) sys.exit(1) sources = [src] elif args.direction: sources = get_enabled_sources(crawler="rss", direction=args.direction) else: sources = get_enabled_sources(crawler="rss") if not sources: log.info("no rss sources to crawl") return for src in sources: run_source(src, dry_run=args.dry_run) time.sleep(RATE_LIMIT_SEC) if __name__ == "__main__": main()