diff --git a/api/vk_parser.py b/api/vk_parser.py index 04f40d4..609c239 100644 --- a/api/vk_parser.py +++ b/api/vk_parser.py @@ -388,13 +388,17 @@ async def _get_or_create_tag(name: str, db: AsyncSession) -> Tag: # ── Post processing ─────────────────────────────────────────────────────────── -async def _process_post(post: dict, category_name: str, db: AsyncSession, strict_filter: bool = True, screen_name: str | None = None) -> bool: +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 False + return "exists" text = post.get("text", "") attachments = post.get("attachments", []) @@ -410,7 +414,7 @@ async def _process_post(post: dict, category_name: str, db: AsyncSession, strict )).scalar_one_or_none() if in_db: # Оригинал придёт сам из своей группы — пропускаем - return False + return "skipped" # Берём содержимое оригинала; комментарий репостера добавляем в начало orig_text = original.get("text", "") orig_atts = original.get("attachments", []) @@ -421,7 +425,7 @@ async def _process_post(post: dict, category_name: str, db: AsyncSession, strict # strict_filter=True: пропускаем пост если нет тегов # strict_filter=False: импортируем всё, теги назначаем если найдены if strict_filter and not tag_names: - return False + return "skipped" # Заранее генерируем UUID и slug — используются в путях MinIO article_id = uuid.uuid4() @@ -545,7 +549,7 @@ async def _process_post(post: dict, category_name: str, db: AsyncSession, strict content = _sanitize("\n".join(content_parts)) if not content.strip(): - return False + return "skipped" title = _early_title title_slug = _early_slug @@ -594,7 +598,7 @@ async def _process_post(post: dict, category_name: str, db: AsyncSession, strict .values(article_id=article_id) ) await db.commit() - return True + return "imported" # ── VK API fetch ────────────────────────────────────────────────────────────── @@ -674,13 +678,16 @@ async def run_import() -> dict: break batch_new = 0 + batch_existing = 0 for post in posts: try: - ok = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name) - if ok: + 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 @@ -691,7 +698,16 @@ async def run_import() -> dict: await asyncio.sleep(1.5) - if batch_new == 0 or len(posts) < batch or not VK_TOKEN: + # Условие остановки: + # — строгий фильтр: новых тегированных постов в батче нет → догнали. + # — без фильтра: останавливаемся только когда ВЕСЬ батч уже в БД, + # иначе пропустили бы более старые нетегированные посты (ранняя остановка). + 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 @@ -749,8 +765,8 @@ async def run_history_import(group_id: str | None = None, since_days: int = 365) hit_cutoff = True continue try: - ok = await _process_post(post, source.group_name, db, source.strict_filter) - if ok: + 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