new web ract

This commit is contained in:
jze9
2026-05-15 03:31:28 +05:00
parent de78624495
commit 2335497226
58 changed files with 7397 additions and 406 deletions

725
api/vk_parser.py Normal file
View File

@@ -0,0 +1,725 @@
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 место",
"лучший студент", "гордост", "поздравляем", "медал"],
"студенческая жизнь": ["студсовет", "студенческ жизн", "общежити", "капустник",
"квн", "студенч", "первокурсник", "посвящение",
"студент год", "актив"],
"профессионалитет": ["профессионалитет", "фп профессионалитет",
"кластер", "федеральн проект"],
# ── Направления колледжей ─────────────────────────────────────────────────
"медицина": ["медицин", "сестринск", "фельдшер", "фармацевт", "лечебн",
"анатоми", "патологи", "здравоохранен", "санитар",
"первая помощ", "реанимац", "клиническ"],
"педагогика": ["педагог", "воспитател", "учител", "начальн класс",
"дошкольн", "детский сад", "логопед", "коррекцион",
"инклюзив", "тьютор"],
"юриспруденция": ["юрист", "юридическ", "правоохранительн", "право",
"законодательств", "суд", "прокурат", "полиц",
"юриспруденц", "правовед"],
"IT и технологии": ["программирован", "it-", "ит-", "информационн технолог",
"цифров", "кибербезопасн", "веб-разраб", "1с",
"компьютерн", "разработчик", "python", "frontend",
"backend", "хакатон", "ворлдскиллс"],
"кулинария и торговля": ["кулинар", "повар", "кондитер", "гастроном", "блюд",
"рецепт", "ресторан", "кафе", "торговл", "продавец",
"мерчандайзинг", "товаровед", "общепит"],
"строительство": ["строительств", "архитектур", "проектирован", "чертёж",
"чертеж", "монтаж", "сварк", "электромонтаж",
"сантехник", "отделочн"],
"техника и механика": ["механик", "двигател", "автомобил", "станок",
"металлообраб", "токар", "слесар", "техническ обслуж",
"ремонт оборудован", "машиностроен", "технолог производ"],
"сельское хозяйство": ["агроном", "агроинженер", "сельск хоз", "животновод",
"агротехник", "землепользован", "агро"],
"экономика и бухучёт": ["бухгалтер", "экономик", "финансов", "налог",
"аудит", "бизнес", "предпринимател", "менеджмент",
"маркетинг", "логистик"],
"социальная работа": ["социальн работ", "соцработник", "психолог",
"реабилитац", "инвалид", "ограниченн возможн",
"социальн помощ", "опека"],
# ── Общественные события ──────────────────────────────────────────────────
"спорт": ["спорт", "соревнован", "олимпиад", "футбол", "баскетбол",
"волейбол", "чемпион", "кубок", "турнир", "атлет",
"фитнес", "зарядк", "ворлдскиллс хайскул"],
"военка": ["военн", "армия", "призыв", "нво", "сво", "оборон",
"патриот", "защитник", "военно-патриот", "стрельб",
"зарниц", "юнармия", "допризывн"],
"воздушная опасность": ["воздушн тревог", "воздушн опасност", "ракет",
"бпла", "дрон", "обстрел", "укрытие", "сирен",
"эвакуац", "бомбоубежищ"],
"общественная деятельность": ["волонтёр", "волонтер", "благотвор", "субботник",
"экологическ акц", "посадк дерев", "уборк территор",
"донорств", "гуманитарн", "помощ фронт"],
"праздники и события": ["праздник", "концерт", "фестиваль", "выставк",
"торжеств", "церемони", "день колледж", "юбилей",
"8 марта", "23 февраля", "новый год", "масленица",
"день знаний", "последний звонок"],
"конкурсы и олимпиады": ["конкурс", "олимпиад", "чемпионат профессий",
"worldskills", "ворлдскиллс", "хакатон", "викторин",
"интеллектуальн", "акселератор"],
}
def _detect_tags(text: str) -> list[str]:
lower = text.lower()
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) -> 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")
object_name = f"vk/images/{uuid.uuid4().hex}.{ext}"
mc = _minio_client()
_ensure_bucket(mc)
mc.put_object(MINIO_BUCKET, object_name, io.BytesIO(data), len(data), content_type=ct)
public_url = f"{MINIO_PUBLIC_URL}/{MINIO_BUCKET}/{object_name}"
if db is not None:
from bd.models import Media
db.add(Media(
filename=object_name.split("/")[-1],
url=public_url,
media_type="image",
size_bytes=len(data),
))
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) -> 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)
object_name = f"vk/videos/{uuid.uuid4().hex}.mp4"
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}/{MINIO_BUCKET}/{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=object_name.split("/")[-1],
url=public_url,
media_type="video",
size_bytes=len(data),
))
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 _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) -> bool:
owner_id = post["owner_id"]
post_id = post["id"]
vk_slug = _vk_slug(owner_id, post_id)
if await _slug_exists(vk_slug, db):
return False
text = post.get("text", "")
attachments = post.get("attachments", [])
ts = post["date"]
# 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
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:
minio_url = await _upload_from_url(url, db=db)
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:
minio_thumb = await _upload_from_url(thumb_url, db=db)
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:
minio_video_url = await _download_vk_video(vid_oid, vid_id, db=db)
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 False
title = _extract_title(text, ts)
title_slug = slugify(title)
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 = _detect_tags(text)
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
source_url = f"https://vk.com/wall{owner_id}_{post_id}"
article = Article(
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()
return True
# ── 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
for post in posts:
try:
ok = await _process_post(post, source.group_name, db)
if ok:
imported += 1; grp_new += 1; batch_new += 1
else:
skipped += 1; grp_skip += 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 batch_new == 0 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:
ok = await _process_post(post, source.group_name, db)
if ok:
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}