Files
news_all_spo/api/vk_parser.py
jze9 45368c64bf fix: VK-импорт без фильтра добирает всю историю в фоне
run_import останавливал группу по batch_new==0, из-за чего для групп
без строгого фильтра (strict_filter=False) фоновый импорт обрывался на
первом батче без новых записей и не доходил до более старых
нетегированных постов.

_process_post теперь возвращает статус imported|exists|skipped, а
условие остановки ветвится: строгий фильтр — прежнее поведение,
без фильтра — стоп только когда весь батч уже в БД.
Заодно run_history_import передаёт screen_name для корректных source_url.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-17 18:47:13 +05:00

796 lines
36 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
import io
import os
import re
import uuid
from datetime import datetime, timezone
import asyncio
import bleach
from bleach.css_sanitizer import CSSSanitizer
import httpx
from minio import Minio
from minio.error import S3Error
from slugify import slugify
from sqlalchemy import select
from sqlalchemy.ext.asyncio import AsyncSession
from bd.database import AsyncSessionLocal
from bd.models import Article, ArticleStatus, Category, Tag, VkSource
VK_API = "https://api.vk.com/method"
VK_VERSION = "5.199"
VK_TOKEN = os.getenv("VK_ACCESS_TOKEN", "")
VK_IMPORT_AS = os.getenv("VK_IMPORT_STATUS", "published") # published | draft
VK_PER_RUN = int(os.getenv("VK_POSTS_PER_RUN", "100")) # макс 100 за запрос
MINIO_ENDPOINT = os.getenv("MINIO_ENDPOINT", "minio:9000")
MINIO_ACCESS_KEY = os.getenv("MINIO_ACCESS_KEY", "minioadmin")
MINIO_SECRET_KEY = os.getenv("MINIO_SECRET_KEY", "minioadmin")
MINIO_BUCKET = os.getenv("MINIO_BUCKET", "news-media")
MINIO_PUBLIC_URL = os.getenv("MINIO_PUBLIC_URL", "http://localhost:9000")
_URL_RE = re.compile(r"(https?://[^\s<>\"']+)")
# Разрешённые HTML-теги и атрибуты после санитизации
_ALLOWED_TAGS = [
"p", "br", "strong", "em", "a", "img",
"ul", "ol", "li", "blockquote",
"div", "iframe", "span", "video",
]
_ALLOWED_ATTRS = {
"a": ["href", "target", "rel"],
"img": ["src", "style", "alt"],
"iframe": ["src", "style", "frameborder", "allowfullscreen"],
"video": ["src", "controls", "style", "poster"],
"div": ["style"],
"span": ["style"],
}
_CSS = CSSSanitizer(allowed_css_properties=[
"max-width", "width", "height", "position", "top", "left",
"padding-bottom", "overflow", "margin", "border-radius",
])
def _sanitize(html: str) -> str:
return bleach.clean(html, tags=_ALLOWED_TAGS, attributes=_ALLOWED_ATTRS,
css_sanitizer=_CSS, strip=True)
# ── Теги по ключевым словам ───────────────────────────────────────────────────
# Ключ = название тега, значение = список подстрок (регистронезависимо)
TAG_KEYWORDS: dict[str, list[str]] = {
"день открытых дверей": ["день открытых дверей", "день открытых",
"приходи в колледж", "познакомиться с колледж",
"двери открыты", "приглашаем на экскурс"],
"общественная деятельность": ["волонтёр", "волонтер", "благотвор", "субботник",
"экологическ акц", "посадк дерев", "уборк территор",
"донорств", "гуманитарн помощ", "добровольч",
"общественн деятельн", "социальн акц", "акция добра"],
"чемпионаты и конкурсы": ["чемпионат",
"олимпиад",
"хакатон",
"всероссийск конкурс", "региональн конкурс",
"областн конкурс", "международн конкурс",
"городск конкурс", "межвузовск конкурс",
"конкурс профессионал", "конкурс мастерств",
"конкурс молодых специалист",
"заняли 1 место", "заняли 2 место", "заняли 3 место",
"1 место среди", "2 место среди", "3 место среди",
"заняли первое место", "заняли второе место", "заняли третье место",
"золото на", "серебро на", "бронзу на",
"завоевали золот", "завоевали серебр", "завоевали бронз"],
"профессионалы": ["профессионалитет", "worldskills", "ворлдскиллс",
"ворлдскиллс хайскул", "хайскул",
"компетенц", "демонстрацион экзамен",
"федеральн проект профессионалите"],
"абилимпикс": ["абилимпикс", "abilympics",
"ограниченн возможн здоровь",
"особенн образовател потребн",
"инклюзивн образован"],
"экзамены": ["расписани экзамен", "подготовк к экзамен",
"экзаменацион сессия", "зачётн неделя", "зачетн неделя",
"промежуточн аттестац", "государственн итог аттестац",
"государственн экзамен", "итоговая аттестац",
"защит диплом", "защита выпускн",
"огэ", "егэ", "гиа"],
}
# Слова-исключения: если они есть в тексте — пост скипается полностью
# (посты о ВОВ/Победе не относятся ни к одной из 6 категорий)
_SKIP_KEYWORDS = [
"великая отечественная война",
"великой отечественной войн",
"день победы",
"9 мая",
"вов ",
"ветеран войн",
"павш за родин",
]
def _detect_tags(text: str) -> list[str]:
lower = text.lower()
# Пропускаем посты о ВОВ/Победе — они не попадают ни в одну категорию
if any(kw in lower for kw in _SKIP_KEYWORDS):
return []
found = []
for tag_name, keywords in TAG_KEYWORDS.items():
if any(kw in lower for kw in keywords):
found.append(tag_name)
return found
# ── MinIO helpers ─────────────────────────────────────────────────────────────
def _minio_client() -> Minio:
return Minio(MINIO_ENDPOINT, access_key=MINIO_ACCESS_KEY, secret_key=MINIO_SECRET_KEY, secure=False)
def _ensure_bucket(mc: Minio):
try:
if not mc.bucket_exists(MINIO_BUCKET):
mc.make_bucket(MINIO_BUCKET)
policy = (
'{"Version":"2012-10-17","Statement":[{"Effect":"Allow","Principal":"*",'
f'"Action":["s3:GetObject"],"Resource":["arn:aws:s3:::{MINIO_BUCKET}/*"]'
'}]}'
)
mc.set_bucket_policy(MINIO_BUCKET, policy)
except S3Error:
pass
async def _upload_from_url(url: str, db: AsyncSession | None = None,
article_dir: str | None = None,
index: int = 1) -> str | None:
try:
async with httpx.AsyncClient(timeout=30, follow_redirects=True) as client:
r = await client.get(url)
if r.status_code != 200:
return None
data = r.content
ct = r.headers.get("content-type", "image/jpeg").split(";")[0].strip()
ext = ct.split("/")[-1].replace("jpeg", "jpg")
if article_dir:
prefix = f"articles/{article_dir}/images"
filename = f"{article_dir}-{index}.{ext}"
else:
prefix = "vk/images"
filename = f"{uuid.uuid4().hex}.{ext}"
object_name = f"{prefix}/{filename}"
mc = _minio_client()
_ensure_bucket(mc)
mc.put_object(MINIO_BUCKET, object_name, io.BytesIO(data), len(data), content_type=ct)
# Системный nginx проксирует /media/ → MinIO bucket, поэтому bucket не дублируем в URL
public_url = f"{MINIO_PUBLIC_URL}/{object_name}"
if db is not None:
from bd.models import Media
db.add(Media(
filename=filename,
url=public_url,
media_type="image",
size_bytes=len(data),
# article_id выставляется после коммита статьи (см. _process_post)
))
return public_url
except Exception as exc:
print(f"[VK] Не удалось загрузить фото {url}: {exc}")
return None
# ── Video download ────────────────────────────────────────────────────────────
async def _fetch_video_info(owner_id: int, video_id: int) -> dict:
"""Получить info о видео через video.get.
Возвращает thumb_url (лучший кадр превью).
files.mp4_* и player возвращаются только если токен имеет scope=video;
без scope они None — в этом случае скачивание идёт через yt-dlp.
"""
if not VK_TOKEN:
return {}
params = {
"videos": f"{owner_id}_{video_id}",
"access_token": VK_TOKEN,
"v": VK_VERSION,
}
try:
async with httpx.AsyncClient(timeout=15) as client:
r = await client.get(f"{VK_API}/video.get", params=params)
data = r.json()
if "error" in data:
print(f"[VK video.get] {data['error'].get('error_msg')}")
return {}
items = data.get("response", {}).get("items", [])
if not items:
return {}
v = items[0]
result = {}
# лучший кадр превью (выше качеством чем в wall.get)
for src_key in ("image", "first_frame"):
frames = v.get(src_key, [])
if frames:
best = max(frames, key=lambda f: f.get("width", 0))
result["thumb_url"] = best.get("url")
break
return result
except Exception as exc:
print(f"[VK] video.get {owner_id}_{video_id}: {exc}")
return {}
def _ytdlp_extract_url(vk_url: str, max_mb: int = 150) -> str | None:
"""Синхронно: через yt-dlp достать прямую mp4-ссылку без скачивания файла."""
try:
import yt_dlp
except ImportError:
return None
max_bytes = max_mb * 1024 * 1024
ydl_opts = {
"format": "best[ext=mp4][filesize<?{}]/best[ext=mp4]/best".format(max_bytes),
"quiet": True,
"no_warnings": True,
}
try:
with yt_dlp.YoutubeDL(ydl_opts) as ydl:
info = ydl.extract_info(vk_url, download=False)
url = info.get("url") if info else None
if url:
print(f"[VK] yt-dlp: найден mp4 для {vk_url}")
return url
except Exception as exc:
print(f"[VK] yt-dlp: не удалось получить URL {vk_url}: {exc}")
return None
async def _download_vk_video(owner_id: int, video_id: int,
db: AsyncSession | None = None,
max_mb: int = 150,
article_dir: str | None = None,
index: int = 1) -> str | None:
"""Скачать VK-видео в MinIO и вернуть публичный URL.
Использует yt-dlp для получения прямой mp4-ссылки.
Стриминг — не держит весь файл в памяти.
"""
vk_url = f"https://vk.com/video{owner_id}_{video_id}"
max_bytes = max_mb * 1024 * 1024
# yt-dlp работает синхронно — запускаем в пуле потоков
loop = asyncio.get_event_loop()
mp4_url = await loop.run_in_executor(None, _ytdlp_extract_url, vk_url, max_mb)
if not mp4_url:
return None
_chrome_ua = (
"Mozilla/5.0 (Windows NT 10.0; Win64; x64) "
"AppleWebKit/537.36 (KHTML, like Gecko) "
"Chrome/124.0.0.0 Safari/537.36"
)
try:
async with httpx.AsyncClient(
timeout=300,
follow_redirects=True,
headers={"User-Agent": _chrome_ua, "Referer": "https://vk.com/"},
) as client:
async with client.stream("GET", mp4_url) as r:
if r.status_code != 200:
print(f"[VK] Видео HTTP {r.status_code}: {mp4_url[:80]}")
return None
ct = r.headers.get("content-type", "video/mp4").split(";")[0].strip()
# Если content-type не видео — значит получили HTML, не файл
if "video" not in ct and "octet" not in ct:
print(f"[VK] Неверный content-type видео: {ct}")
return None
declared = int(r.headers.get("content-length", 0))
if declared and declared > max_bytes:
print(f"[VK] Видео слишком большое: {declared/1024/1024:.0f}MB > {max_mb}MB, пропускаем")
return None
chunks, total = [], 0
async for chunk in r.aiter_bytes(1024 * 256):
total += len(chunk)
if total > max_bytes:
print(f"[VK] Видео превысило {max_mb}MB при скачивании, пропускаем")
return None
chunks.append(chunk)
data = b"".join(chunks)
if article_dir:
prefix = f"articles/{article_dir}/videos"
filename = f"{article_dir}-{index}.mp4"
else:
prefix = "vk/videos"
filename = f"{uuid.uuid4().hex}.mp4"
object_name = f"{prefix}/{filename}"
mc = _minio_client()
_ensure_bucket(mc)
mc.put_object(MINIO_BUCKET, object_name, io.BytesIO(data), len(data), content_type="video/mp4")
public_url = f"{MINIO_PUBLIC_URL}/{object_name}"
print(f"[VK] Видео загружено в MinIO: {object_name} ({total/1024/1024:.1f}MB)")
if db is not None:
from bd.models import Media
db.add(Media(
filename=filename,
url=public_url,
media_type="video",
size_bytes=len(data),
# article_id выставляется после коммита статьи (см. _process_post)
))
return public_url
except Exception as exc:
print(f"[VK] Ошибка скачивания видео {vk_url}: {exc}")
return None
# ── VK content helpers ────────────────────────────────────────────────────────
def _best_photo(sizes: list) -> str | None:
for t in ("w", "z", "y", "x", "m", "s"):
for s in sizes:
if s.get("type") == t:
return s["url"]
return sizes[-1]["url"] if sizes else None
def _text_to_html(text: str) -> str:
parts = []
for line in text.split("\n"):
line = line.strip()
if not line:
continue
line = _URL_RE.sub(r'<a href="\1" target="_blank">\1</a>', line)
parts.append(f"<p>{line}</p>")
return "\n".join(parts)
def _extract_title(text: str, ts: int) -> str:
for line in text.split("\n"):
line = line.strip()
if line:
return line[:200]
return "ВКонтакте " + datetime.utcfromtimestamp(ts).strftime("%d.%m.%Y")
def _vk_slug(owner_id: int, post_id: int) -> str:
return f"vk-{abs(owner_id)}-{post_id}"
# ── DB helpers ────────────────────────────────────────────────────────────────
async def _slug_exists(slug: str, db: AsyncSession) -> bool:
r = await db.execute(select(Article.id).where(Article.slug == slug))
return r.scalar_one_or_none() is not None
async def _vk_post_exists(owner_id: int, post_id: int, db: AsyncSession) -> bool:
pattern = f"%wall{owner_id}_{post_id}%"
r = await db.execute(select(Article.id).where(Article.source_url.like(pattern)))
return r.scalar_one_or_none() is not None
async def _get_or_create_category(name: str, db: AsyncSession) -> Category:
s = slugify(name)
cat = (await db.execute(select(Category).where(Category.slug == s))).scalar_one_or_none()
if not cat:
cat = Category(name=name, slug=s)
db.add(cat)
await db.flush()
return cat
async def _get_or_create_tag(name: str, db: AsyncSession) -> Tag:
s = slugify(name)
tag = (await db.execute(select(Tag).where(Tag.slug == s))).scalar_one_or_none()
if not tag:
tag = Tag(name=name, slug=s)
db.add(tag)
await db.flush()
return tag
# ── Post processing ───────────────────────────────────────────────────────────
async def _process_post(post: dict, category_name: str, db: AsyncSession, strict_filter: bool = True, screen_name: str | None = None) -> str:
"""Обрабатывает один VK-пост. Возвращает статус:
"imported" — пост сохранён; "exists" — уже есть в БД (догнали историю);
"skipped" — пропущен по другой причине (репост известной группы / фильтр / пустой).
"""
owner_id = post["owner_id"]
post_id = post["id"]
vk_slug = _vk_slug(owner_id, post_id)
if await _vk_post_exists(owner_id, post_id, db):
return "exists"
text = post.get("text", "")
attachments = post.get("attachments", [])
ts = post["date"]
# Обработка репостов (copy_history)
copy_history = post.get("copy_history", [])
if copy_history:
original = copy_history[0]
orig_gid = str(abs(original.get("owner_id", 0)))
in_db = (await db.execute(
select(VkSource.id).where(VkSource.group_id == orig_gid)
)).scalar_one_or_none()
if in_db:
# Оригинал придёт сам из своей группы — пропускаем
return "skipped"
# Берём содержимое оригинала; комментарий репостера добавляем в начало
orig_text = original.get("text", "")
orig_atts = original.get("attachments", [])
text = (text + "\n" + orig_text).strip() if text else orig_text
attachments = orig_atts if orig_atts else attachments
tag_names = _detect_tags(text)
# strict_filter=True: пропускаем пост если нет тегов
# strict_filter=False: импортируем всё, теги назначаем если найдены
if strict_filter and not tag_names:
return "skipped"
# Заранее генерируем UUID и slug — используются в путях MinIO
article_id = uuid.uuid4()
# VK Clips и некоторые видео-посты: post.text пустой, описание лежит в video.description
if not text:
for att in attachments:
if att.get("type") == "video":
desc = att["video"].get("description", "").strip()
if desc:
text = desc
break
# Генерируем slug до загрузки медиа: используем его как имя директории и файлов
_early_title = _extract_title(text, ts)
_early_slug = (slugify(_early_title) or vk_slug)[:60] # макс 60 символов
# Уникальный префикс: slug + 8 символов UUID (защита от коллизий одинаковых заголовков)
article_dir = f"{_early_slug}-{str(article_id)[:8]}"
img_idx = 0 # счётчик изображений для нумерации файлов
vid_idx = 0 # счётчик видео
cover_url = None
content_parts = []
if text:
content_parts.append(_text_to_html(text))
for att in attachments:
att_type = att.get("type")
if att_type == "photo":
url = _best_photo(att["photo"].get("sizes", []))
if url:
img_idx += 1
minio_url = await _upload_from_url(url, db=db, article_dir=article_dir, index=img_idx)
final_url = minio_url or url # fallback на оригинальный URL если MinIO недоступен
if cover_url is None:
cover_url = final_url
else:
content_parts.append(
f'<p><img src="{final_url}" style="max-width:100%;border-radius:8px;"></p>'
)
elif att_type == "video":
vid = att["video"]
player = vid.get("player")
vtitle = vid.get("title", "Видео")
vid_id = vid.get("id")
vid_oid = vid.get("owner_id")
vk_video_url = (
f"https://vk.com/video{vid_oid}_{vid_id}"
if vid_id and vid_oid else None
)
# Кадр превью из wall.get (запасной)
thumb_url = None
for src_key in ("image", "first_frame"):
frames = vid.get(src_key, [])
if frames:
best = max(frames, key=lambda f: f.get("width", 0))
thumb_url = best.get("url")
break
# Улучшенный кадр превью из video.get (+ может вернуть player при scope=video)
if vid_id and vid_oid:
vinfo = await _fetch_video_info(vid_oid, vid_id)
if vinfo.get("thumb_url"):
thumb_url = vinfo["thumb_url"]
# Всегда скачиваем превью в MinIO
minio_thumb = None
if thumb_url:
img_idx += 1
minio_thumb = await _upload_from_url(thumb_url, db=db, article_dir=article_dir, index=img_idx)
final_thumb = minio_thumb or thumb_url
if cover_url is None:
cover_url = final_thumb
else:
final_thumb = None
# Скачиваем само видео через yt-dlp → MinIO
minio_video_url = None
if vid_id and vid_oid:
vid_idx += 1
minio_video_url = await _download_vk_video(vid_oid, vid_id, db=db, article_dir=article_dir, index=vid_idx)
if minio_video_url:
# 1. Видео лежит у нас — нативный <video>
poster_attr = f' poster="{final_thumb}"' if final_thumb else ""
content_parts.append(
f'<div style="margin:1em 0">'
f'<video src="{minio_video_url}"{poster_attr} controls '
f'style="max-width:100%;border-radius:8px;display:block;">'
f'</video>'
f'<p><em>{vtitle}</em></p></div>'
)
elif vk_video_url:
# 2. Fallback: превью + ссылка на ВК (yt-dlp не смог)
if final_thumb:
content_parts.append(
f'<p><a href="{vk_video_url}" target="_blank" rel="noopener">'
f'<img src="{final_thumb}" alt="{vtitle}" '
f'style="max-width:100%;border-radius:8px;display:block;">'
f'</a></p>'
f'<p>▶ <a href="{vk_video_url}" target="_blank" rel="noopener">'
f'{vtitle}</a></p>'
)
else:
content_parts.append(
f'<p>▶ <a href="{vk_video_url}" target="_blank" rel="noopener">'
f'{vtitle}</a></p>'
)
elif att_type == "link":
href = att["link"].get("url", "")
label = att["link"].get("title") or href
if href:
content_parts.append(f'<p>🔗 <a href="{href}" target="_blank">{label}</a></p>')
content = _sanitize("\n".join(content_parts))
if not content.strip():
return "skipped"
title = _early_title
title_slug = _early_slug
final_slug = title_slug if not await _slug_exists(title_slug, db) else vk_slug
# Категория = название группы
category = await _get_or_create_category(category_name, db)
# Теги уже определены выше (tag_names)
tags = [await _get_or_create_tag(n, db) for n in tag_names]
status = ArticleStatus.published if VK_IMPORT_AS == "published" else ArticleStatus.draft
# Храним naive UTC — колонка TIMESTAMP WITHOUT TIME ZONE
created_at = datetime.utcfromtimestamp(ts)
published_at = created_at if status == ArticleStatus.published else None
if screen_name:
source_url = f"https://vk.com/{screen_name}?w=wall{owner_id}_{post_id}"
else:
source_url = f"https://vk.com/wall{owner_id}_{post_id}"
article = Article(
id=article_id,
title=title,
slug=final_slug,
content=content,
excerpt=text[:300].strip(),
cover_url=cover_url,
source_url=source_url,
status=status,
published_at=published_at,
created_at=created_at,
updated_at=created_at,
category_id=category.id,
tags=tags,
)
db.add(article)
await db.commit()
# Связываем загруженные медиафайлы с этой статьёй по URL-префиксу
from sqlalchemy import update as sa_update
from bd.models import Media as _Media
await db.execute(
sa_update(_Media)
.where(_Media.url.contains(f"/articles/{article_dir}/"))
.values(article_id=article_id)
)
await db.commit()
return "imported"
# ── VK API fetch ──────────────────────────────────────────────────────────────
async def _fetch_posts(group_id: str, count: int) -> list:
# Если токен есть — официальный API, иначе — HTML-парсинг
if VK_TOKEN:
return await _fetch_posts_api(group_id, count)
return await _fetch_posts_html(group_id, count)
async def _fetch_posts_api(group_id: str, count: int, offset: int = 0) -> list:
owner_id = f"-{group_id.lstrip('-')}"
params = {
"owner_id": owner_id,
"count": count,
"offset": offset,
"filter": "owner",
"access_token": VK_TOKEN,
"v": VK_VERSION,
}
try:
async with httpx.AsyncClient(timeout=15) as client:
r = await client.get(f"{VK_API}/wall.get", params=params)
data = r.json()
if "error" in data:
print(f"[VK API] Ошибка группы {group_id}: {data['error'].get('error_msg')}")
return []
items = data.get("response", {}).get("items", [])
print(f"[VK API] Группа {group_id} offset={offset}: {len(items)} постов")
return items
except Exception as exc:
print(f"[VK API] Сетевая ошибка группы {group_id}: {exc}")
return []
async def _fetch_posts_html(group_id: str, count: int) -> list:
from vk_scraper import fetch_group_posts_html
posts = await fetch_group_posts_html(group_id, count)
print(f"[VK scraper] Группа {group_id}: {len(posts)} постов")
return posts
# ── Public entry point ────────────────────────────────────────────────────────
async def run_import() -> dict:
"""Импортирует только новые посты: пагинирует каждую группу и останавливается,
когда весь батч состоит из уже сохранённых постов."""
mode = "API" if VK_TOKEN else "HTML-scraper"
imported = skipped = errors = 0
batch = 100
async with AsyncSessionLocal() as db:
sources = (await db.execute(
select(VkSource).where(VkSource.enabled == True)
)).scalars().all()
if not sources:
print("[VK] Нет активных источников")
return {"ok": False, "reason": "Нет активных VK-источников"}
print(f"[VK] ▶ Импорт новых постов | групп: {len(sources)} | режим: {mode}")
for i, source in enumerate(sources, 1):
src_offset = 0
grp_new = grp_skip = grp_err = 0
print(f"[VK] [{i}/{len(sources)}] Группа «{source.group_name}» (id={source.group_id})")
while True:
if VK_TOKEN:
posts = await _fetch_posts_api(source.group_id, batch, src_offset)
else:
posts = await _fetch_posts_html(source.group_id, batch)
if not posts:
print(f"[VK] offset={src_offset}: постов нет, завершаем группу")
break
batch_new = 0
batch_existing = 0
for post in posts:
try:
st = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name)
if st == "imported":
imported += 1; grp_new += 1; batch_new += 1
else:
skipped += 1; grp_skip += 1
if st == "exists":
batch_existing += 1
except Exception as exc:
print(f"[VK] ✗ Пост {post.get('id')}: {exc}")
errors += 1; grp_err += 1
await db.rollback()
print(f"[VK] offset={src_offset}: получено {len(posts)}, "
f"новых +{batch_new}, пропущено {len(posts)-batch_new}")
await asyncio.sleep(1.5)
# Условие остановки:
# — строгий фильтр: новых тегированных постов в батче нет → догнали.
# — без фильтра: останавливаемся только когда ВЕСЬ батч уже в БД,
# иначе пропустили бы более старые нетегированные посты (ранняя остановка).
if source.strict_filter:
caught_up = batch_new == 0
else:
caught_up = batch_existing == len(posts)
if caught_up or len(posts) < batch or not VK_TOKEN:
print(f"[VK] Догнали до уже импортированных, останавливаем группу")
break
src_offset += batch
print(f"[VK] ✓ «{source.group_name}»: +{grp_new} новых, {grp_skip} пропущено, {grp_err} ошибок")
source.last_run = datetime.utcnow()
await db.commit()
print(f"[VK] ■ Импорт завершён: +{imported} новых, {skipped} уже было, {errors} ошибок")
return {"ok": True, "imported": imported, "skipped": skipped, "errors": errors, "mode": mode}
async def run_history_import(group_id: str | None = None, since_days: int = 365) -> dict:
"""Пагинированный импорт постов за последние since_days дней."""
from datetime import timedelta
cutoff = (datetime.now(timezone.utc) - timedelta(days=since_days)).timestamp()
cutoff_str = datetime.utcfromtimestamp(cutoff).strftime("%d.%m.%Y")
mode = "API" if VK_TOKEN else "HTML-scraper"
imported = skipped = errors = 0
batch = 100
async with AsyncSessionLocal() as db:
q = select(VkSource).where(VkSource.enabled == True)
if group_id:
q = q.where(VkSource.group_id == group_id)
sources = (await db.execute(q)).scalars().all()
if not sources:
print("[VK] Источники не найдены")
return {"ok": False, "reason": "Источники не найдены"}
print(f"[VK] ▶ Исторический импорт | с {cutoff_str} | групп: {len(sources)} | режим: {mode}")
for i, source in enumerate(sources, 1):
src_offset = 0
grp_new = grp_skip = grp_err = 0
print(f"[VK] [{i}/{len(sources)}] Группа «{source.group_name}» (id={source.group_id})")
while True:
if VK_TOKEN:
posts = await _fetch_posts_api(source.group_id, batch, src_offset)
else:
posts = await _fetch_posts_html(source.group_id, batch)
if not posts:
print(f"[VK] offset={src_offset}: постов нет, завершаем группу")
break
hit_cutoff = False
batch_new = 0
for post in posts:
post_date = datetime.utcfromtimestamp(post.get("date", 0)).strftime("%d.%m.%Y")
if post.get("date", 0) < cutoff:
hit_cutoff = True
continue
try:
st = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name)
if st == "imported":
imported += 1; grp_new += 1; batch_new += 1
else:
skipped += 1; grp_skip += 1
except Exception as exc:
print(f"[VK] ✗ Пост {post.get('id')} ({post_date}): {exc}")
errors += 1; grp_err += 1
await db.rollback()
oldest = datetime.utcfromtimestamp(posts[-1].get("date", 0)).strftime("%d.%m.%Y")
print(f"[VK] offset={src_offset}: получено {len(posts)}, "
f"новых +{batch_new}, пропущено {grp_skip}, самый старый: {oldest}"
+ (" ← достигли года" if hit_cutoff else ""))
await asyncio.sleep(2.0)
if hit_cutoff or len(posts) < batch or not VK_TOKEN:
break
src_offset += batch
print(f"[VK] ✓ «{source.group_name}»: +{grp_new} новых, {grp_skip} уже было, {grp_err} ошибок")
source.last_run = datetime.utcnow()
await db.commit()
print(f"[VK] ■ Исторический импорт завершён: +{imported} новых, {skipped} уже было, {errors} ошибок")
return {"ok": True, "imported": imported, "skipped": skipped, "errors": errors, "mode": mode}