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>
This commit is contained in:
@@ -388,13 +388,17 @@ async def _get_or_create_tag(name: str, db: AsyncSession) -> Tag:
|
|||||||
|
|
||||||
# ── Post processing ───────────────────────────────────────────────────────────
|
# ── 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"]
|
owner_id = post["owner_id"]
|
||||||
post_id = post["id"]
|
post_id = post["id"]
|
||||||
vk_slug = _vk_slug(owner_id, post_id)
|
vk_slug = _vk_slug(owner_id, post_id)
|
||||||
|
|
||||||
if await _vk_post_exists(owner_id, post_id, db):
|
if await _vk_post_exists(owner_id, post_id, db):
|
||||||
return False
|
return "exists"
|
||||||
|
|
||||||
text = post.get("text", "")
|
text = post.get("text", "")
|
||||||
attachments = post.get("attachments", [])
|
attachments = post.get("attachments", [])
|
||||||
@@ -410,7 +414,7 @@ async def _process_post(post: dict, category_name: str, db: AsyncSession, strict
|
|||||||
)).scalar_one_or_none()
|
)).scalar_one_or_none()
|
||||||
if in_db:
|
if in_db:
|
||||||
# Оригинал придёт сам из своей группы — пропускаем
|
# Оригинал придёт сам из своей группы — пропускаем
|
||||||
return False
|
return "skipped"
|
||||||
# Берём содержимое оригинала; комментарий репостера добавляем в начало
|
# Берём содержимое оригинала; комментарий репостера добавляем в начало
|
||||||
orig_text = original.get("text", "")
|
orig_text = original.get("text", "")
|
||||||
orig_atts = original.get("attachments", [])
|
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=True: пропускаем пост если нет тегов
|
||||||
# strict_filter=False: импортируем всё, теги назначаем если найдены
|
# strict_filter=False: импортируем всё, теги назначаем если найдены
|
||||||
if strict_filter and not tag_names:
|
if strict_filter and not tag_names:
|
||||||
return False
|
return "skipped"
|
||||||
|
|
||||||
# Заранее генерируем UUID и slug — используются в путях MinIO
|
# Заранее генерируем UUID и slug — используются в путях MinIO
|
||||||
article_id = uuid.uuid4()
|
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))
|
content = _sanitize("\n".join(content_parts))
|
||||||
if not content.strip():
|
if not content.strip():
|
||||||
return False
|
return "skipped"
|
||||||
|
|
||||||
title = _early_title
|
title = _early_title
|
||||||
title_slug = _early_slug
|
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)
|
.values(article_id=article_id)
|
||||||
)
|
)
|
||||||
await db.commit()
|
await db.commit()
|
||||||
return True
|
return "imported"
|
||||||
|
|
||||||
|
|
||||||
# ── VK API fetch ──────────────────────────────────────────────────────────────
|
# ── VK API fetch ──────────────────────────────────────────────────────────────
|
||||||
@@ -674,13 +678,16 @@ async def run_import() -> dict:
|
|||||||
break
|
break
|
||||||
|
|
||||||
batch_new = 0
|
batch_new = 0
|
||||||
|
batch_existing = 0
|
||||||
for post in posts:
|
for post in posts:
|
||||||
try:
|
try:
|
||||||
ok = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name)
|
st = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name)
|
||||||
if ok:
|
if st == "imported":
|
||||||
imported += 1; grp_new += 1; batch_new += 1
|
imported += 1; grp_new += 1; batch_new += 1
|
||||||
else:
|
else:
|
||||||
skipped += 1; grp_skip += 1
|
skipped += 1; grp_skip += 1
|
||||||
|
if st == "exists":
|
||||||
|
batch_existing += 1
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
print(f"[VK] ✗ Пост {post.get('id')}: {exc}")
|
print(f"[VK] ✗ Пост {post.get('id')}: {exc}")
|
||||||
errors += 1; grp_err += 1
|
errors += 1; grp_err += 1
|
||||||
@@ -691,7 +698,16 @@ async def run_import() -> dict:
|
|||||||
|
|
||||||
await asyncio.sleep(1.5)
|
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] Догнали до уже импортированных, останавливаем группу")
|
print(f"[VK] Догнали до уже импортированных, останавливаем группу")
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -749,8 +765,8 @@ async def run_history_import(group_id: str | None = None, since_days: int = 365)
|
|||||||
hit_cutoff = True
|
hit_cutoff = True
|
||||||
continue
|
continue
|
||||||
try:
|
try:
|
||||||
ok = await _process_post(post, source.group_name, db, source.strict_filter)
|
st = await _process_post(post, source.group_name, db, source.strict_filter, screen_name=source.screen_name)
|
||||||
if ok:
|
if st == "imported":
|
||||||
imported += 1; grp_new += 1; batch_new += 1
|
imported += 1; grp_new += 1; batch_new += 1
|
||||||
else:
|
else:
|
||||||
skipped += 1; grp_skip += 1
|
skipped += 1; grp_skip += 1
|
||||||
|
|||||||
Reference in New Issue
Block a user