diff --git a/db/migrations/046_insta_parser_settings.sql b/db/migrations/046_insta_parser_settings.sql new file mode 100644 index 0000000..96b81b7 --- /dev/null +++ b/db/migrations/046_insta_parser_settings.sql @@ -0,0 +1,10 @@ +INSERT INTO app_settings (key, label, value_type, value) +VALUES +('insta_login', 'Instagram Login', 'string', ''), +('insta_password', 'Instagram Password', 'secret', ''), +('insta_proxy_url', 'Instagram Proxy URL (http://user:pass@ip:port)', 'string', ''), +('insta_fetch_count', 'Instagram: количество проверяемых последних постов', 'int', '5'), +('insta_delay_base_minutes', 'Instagram: пауза между аккаунтами (минут)', 'int', '35'), +('insta_delay_random_minutes', 'Instagram: разброс паузы (± минут)', 'int', '5'), +('insta_cooldown_hours', 'Instagram: отлежка при лимитах (часов)', 'int', '12') +ON CONFLICT (key) DO NOTHING; diff --git a/requirements.txt b/requirements.txt index 7c83571..4a9fb58 100644 --- a/requirements.txt +++ b/requirements.txt @@ -10,3 +10,4 @@ Pillow==11.3.0 python-multipart==0.0.20 uvicorn[standard]==0.35.0 yt-dlp==2026.6.9 +instagrapi==2.1.2 diff --git a/src/vk_parser_app/admin.py b/src/vk_parser_app/admin.py index b0163f7..151717e 100644 --- a/src/vk_parser_app/admin.py +++ b/src/vk_parser_app/admin.py @@ -22,7 +22,7 @@ from fastapi.templating import Jinja2Templates from loguru import logger from .config import settings -from .constants import PLATFORM_VK +from .constants import PLATFORM_VK, PLATFORM_INSTAGRAM from .db import fetch_int_setting, fetch_setting, get_pool from .security import hash_password, new_token, token_hash, verify_password from .text_utils import build_publication_text, normalize_hash_tag, parse_categories @@ -37,6 +37,7 @@ from .workers.tg_poster import TelegramPoster from .workers.tg_reactor import TelegramReactor from .workers.vk_poster import VKPoster from .workers.vk_storage_uploader import TelegramStorageUploader +from .workers.insta_parser import InstaParserWorker COOKIE_NAME = "vk_parser_admin" VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier" @@ -1811,6 +1812,7 @@ async def startup() -> None: ("tg-reactor", TelegramReactor()), ("vk-poster", VKPoster()), ("vk-storage-uploader", TelegramStorageUploader()), + ("insta-parser", InstaParserWorker()), ] for name, worker in workers: asyncio.create_task(start_worker_task(worker, name)) @@ -2316,6 +2318,16 @@ async def sources_bulk_create( if not isinstance(item, dict) or not item.get("ok"): continue try: + url_val = str(item.get("url") or "").strip() + external_id_val = str(item.get("external_id") or "").strip() + platform_val = PLATFORM_VK + if "instagram.com" in url_val: + platform_val = PLATFORM_INSTAGRAM + if not external_id_val: + m = re.search(r"instagram\.com/([^/?#]+)", url_val) + if m: + external_id_val = m.group(1) + row = await pool.fetchrow( """ INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, created_by) @@ -2323,11 +2335,11 @@ async def sources_bulk_create( ON CONFLICT DO NOTHING RETURNING id """, - PLATFORM_VK, - str(item.get("name") or item.get("external_id") or "").strip(), - normalize_hash_tag(str(item.get("tag") or item.get("external_id") or ""), str(item.get("external_id") or "source")), - str(item.get("url") or "").strip(), - str(item.get("external_id") or "").strip(), + platform_val, + str(item.get("name") or external_id_val or "").strip(), + normalize_hash_tag(str(item.get("tag") or external_id_val or ""), external_id_val or "source"), + url_val, + external_id_val, item.get("external_owner_id"), bool(item.get("active", True)), user["id"], @@ -2367,10 +2379,21 @@ async def source_create( return redirect("/login") require_csrf(user, csrf_token) platform = platform.strip().lower() or PLATFORM_VK - external_id = normalize_vk_source(url) if platform == PLATFORM_VK else "" + resolved_url = url.strip() + + if "instagram.com" in resolved_url and platform == PLATFORM_VK: + platform = PLATFORM_INSTAGRAM + + if platform == PLATFORM_VK: + external_id = normalize_vk_source(resolved_url) + elif platform == PLATFORM_INSTAGRAM: + m = re.search(r"instagram\.com/([^/?#]+)", resolved_url) + external_id = m.group(1) if m else "" + else: + external_id = "" + external_owner_id = None resolved_name = name.strip() - resolved_url = url.strip() status_value = "new" status_msg = None if platform == PLATFORM_VK and external_id: diff --git a/src/vk_parser_app/constants.py b/src/vk_parser_app/constants.py index 35d313f..612aa81 100644 --- a/src/vk_parser_app/constants.py +++ b/src/vk_parser_app/constants.py @@ -1,4 +1,5 @@ PLATFORM_VK = "vk" +PLATFORM_INSTAGRAM = "instagram" SOURCE_STATUS_NEW = "new" SOURCE_STATUS_OK = "ok" @@ -38,3 +39,4 @@ WORKER_VK_POSTER = "vk-poster" WORKER_MAX_POSTER = "max-poster" WORKER_SITE_POSTER = "site-poster" WORKER_DAILY_REPORT = "daily-report" +WORKER_INSTA_PARSER = "insta-parser" diff --git a/src/vk_parser_app/workers/ai_alerts.py b/src/vk_parser_app/workers/ai_alerts.py index b687b01..46209e0 100644 --- a/src/vk_parser_app/workers/ai_alerts.py +++ b/src/vk_parser_app/workers/ai_alerts.py @@ -69,3 +69,37 @@ async def send_ai_worker_error_alert(worker_name: str, model: str, post_ids: lis await asyncio.sleep(float(exc.retry_after) + 1) finally: await bot.session.close() + + +async def send_parser_error_alert(parser_name: str, source_name: str, error: str) -> None: + token = ( + str(await fetch_setting("daily_report_bot_token", "") or "").strip() + or str(await fetch_setting("tg_poster_bot_token", "") or "").strip() + or settings.tg_bot_token + ) + recipients = parse_recipients(await fetch_setting("daily_report_recipient_ids", [442509142])) + if not token or not recipients: + logger.warning("Parser alert skipped: token or recipients are empty") + return + + text = ( + f"🚨 CRITICAL PARSER ERROR\n" + f"Parser: {parser_name}\n" + f"Source: {source_name}\n\n" + f"Error: {error[:1000]}" + ) + bot = Bot(token=token) + try: + for recipient_id in recipients: + while True: + try: + await bot.send_message(recipient_id, text, disable_web_page_preview=True, parse_mode="HTML") + break + except TelegramRetryAfter as exc: + await asyncio.sleep(float(exc.retry_after) + 1) + except Exception as e: + logger.error("Failed to send parser alert to {}: {}", recipient_id, e) + break + finally: + await bot.session.close() + diff --git a/src/vk_parser_app/workers/insta_parser.py b/src/vk_parser_app/workers/insta_parser.py new file mode 100644 index 0000000..61170a0 --- /dev/null +++ b/src/vk_parser_app/workers/insta_parser.py @@ -0,0 +1,448 @@ +from __future__ import annotations + +import asyncio +import hashlib +import json +import random +import time +from datetime import datetime, timedelta, timezone +from pathlib import Path + +from instagrapi import Client +from instagrapi.exceptions import ( + ChallengeRequired, + PleaseWaitFewMinutes, + LoginRequired, + UserNotFound, + PrivateAccount +) +from loguru import logger + +from ..config import settings +from ..constants import ( + PLATFORM_INSTAGRAM, + POST_STATUS_TRASH, + POST_STATUS_MEDIA_PENDING, + WORKER_INSTA_PARSER, +) +from ..db import fetch_setting, fetch_int_setting, get_pool +from ..heartbeat import HeartbeatReporter +from .ai_alerts import send_parser_error_alert + + +def make_content_hash(text: str) -> str: + return hashlib.sha256((text or "").strip().lower().encode()).hexdigest() + + +class InstaParserWorker: + def __init__(self): + self.pool = None + self.heartbeat_interval_sec = 60 + self.heartbeat = HeartbeatReporter("insta-parser", 60) + self.heartbeat_task: asyncio.Task | None = None + self.client: Client | None = None + self.session_file = Path("insta_session.json") + + async def init(self): + self.pool = await get_pool() + + async def get_active_source(self) -> dict | None: + async with self.pool.acquire() as conn: + row = await conn.fetchrow( + """ + SELECT id, external_id, name, last_checked_at, last_parsed_at + FROM sources + WHERE platform = $1 AND active = TRUE + ORDER BY last_checked_at NULLS FIRST + LIMIT 1 + """, + PLATFORM_INSTAGRAM + ) + return dict(row) if row else None + + async def get_known_hashes(self, hashes: list[str]) -> set[str]: + if not hashes: + return set() + async with self.pool.acquire() as conn: + rows = await conn.fetch( + """ + SELECT content_hash + FROM posts + WHERE content_hash = ANY($1::text[]) + """, + hashes, + ) + return {str(r["content_hash"]) for r in rows} + + async def get_known_post_ids(self, source_id: int, external_post_ids: list[str]) -> set[str]: + if not external_post_ids: + return set() + async with self.pool.acquire() as conn: + try: + numeric_ids = [int(pk) for pk in external_post_ids] + rows = await conn.fetch( + """ + SELECT vk_post_id + FROM posts + WHERE source_id = $1 AND vk_post_id = ANY($2::bigint[]) + """, + source_id, + numeric_ids, + ) + return {str(r["vk_post_id"]) for r in rows} + except ValueError: + return set() + + async def deactivate_source(self, source_id: int, reason: str) -> None: + async with self.pool.acquire() as conn: + await conn.execute( + """ + UPDATE sources + SET active = FALSE, + status = 'error', + status_msg = $1, + last_checked_at = NOW(), + updated_at = NOW() + WHERE id = $2 + """, + reason[:1000], + source_id, + ) + + async def mark_source_ok(self, source_id: int, last_parsed_at: datetime | None = None) -> None: + async with self.pool.acquire() as conn: + await conn.execute( + """ + UPDATE sources + SET status = 'ok', + status_msg = NULL, + last_checked_at = NOW(), + last_parsed_at = COALESCE($2, last_parsed_at), + updated_at = NOW() + WHERE id = $1 + """, + source_id, + last_parsed_at, + ) + + async def save_cooldown(self, hours: int) -> None: + until = datetime.now(tz=timezone.utc) + timedelta(hours=hours) + async with self.pool.acquire() as conn: + await conn.execute( + """ + INSERT INTO app_settings (key, value, value_type) + VALUES ('insta_cooldown_until', $1, 'string') + ON CONFLICT (key) DO UPDATE SET value = $1 + """, + until.isoformat() + ) + + async def get_cooldown_until(self) -> datetime | None: + async with self.pool.acquire() as conn: + val = await conn.fetchval("SELECT value FROM app_settings WHERE key = 'insta_cooldown_until'") + if val: + try: + dt = datetime.fromisoformat(val) + if dt.tzinfo is None: + dt = dt.replace(tzinfo=timezone.utc) + return dt + except ValueError: + pass + return None + + async def save_post_and_media(self, source_id: int, post: dict, status: str, skip_reason: str | None) -> int | None: + text = (post.get("caption_text") or "").strip() + content_hash = make_content_hash(text) + posted_at = post.get("taken_at") + if not posted_at: + posted_at = datetime.now(tz=timezone.utc) + elif posted_at.tzinfo is None: + posted_at = posted_at.replace(tzinfo=timezone.utc) + + pk = int(post["pk"]) + + photos = [] + videos = [] + + media_type = post.get("media_type") + if media_type == 1: + url = post.get("thumbnail_url") + if url: photos.append(url) + elif media_type == 2: + url = post.get("video_url") + if url: videos.append(url) + elif media_type == 8: + for resource in post.get("resources", []): + if resource.get("media_type") == 1: + photos.append(resource.get("thumbnail_url")) + elif resource.get("media_type") == 2: + videos.append(resource.get("video_url")) + + async with self.pool.acquire() as conn: + async with conn.transaction(): + post_id = await conn.fetchval( + """ + INSERT INTO posts ( + source_id, vk_post_id, vk_owner_id, posted_at, + raw_text, raw_json, content_hash, + status, skip_reason + ) + VALUES ($1,$2,$3,$4,$5,$6::jsonb,$7,$8::post_status,$9) + ON CONFLICT (source_id, vk_post_id) DO NOTHING + RETURNING id + """, + source_id, + pk, + int(post.get("user", {}).get("pk") or 0), + posted_at, + text, + json.dumps(post, default=str, ensure_ascii=False), + content_hash, + status, + skip_reason, + ) + if post_id is None: + return None + + for p_url in photos: + if not p_url: continue + await conn.execute( + """ + INSERT INTO post_media (post_id, media_type, vk_url) + VALUES ($1, 'photo', $2) + """, + post_id, + str(p_url), + ) + + for v_url in videos: + if not v_url: continue + await conn.execute( + """ + INSERT INTO post_media (post_id, media_type, vk_url) + VALUES ($1, 'video', $2) + """, + post_id, + str(v_url), + ) + + if status == POST_STATUS_MEDIA_PENDING: + await conn.execute( + """ + INSERT INTO jobs (job_type, entity_type, entity_id, payload, status) + VALUES ('media.process', 'post', $1, '{}'::jsonb, 'PENDING') + ON CONFLICT DO NOTHING + """, + post_id, + ) + return post_id + + def setup_client(self, login: str, password: str, proxy: str | None) -> bool: + if self.client: + return True + + self.client = Client() + if proxy: + self.client.set_proxy(proxy) + + if self.session_file.exists(): + try: + self.client.load_settings(self.session_file) + self.client.get_timeline_feed() + logger.info("Instagram session loaded successfully") + return True + except Exception as e: + logger.warning(f"Session invalid, will re-login: {e}") + self.session_file.unlink(missing_ok=True) + + logger.info(f"Logging into Instagram as {login}...") + try: + self.client.login(login, password) + self.client.dump_settings(self.session_file) + logger.info("Instagram login successful, session saved") + return True + except ChallengeRequired as e: + logger.error(f"Challenge Required during login: {e}") + raise + except Exception as e: + logger.error(f"Failed to login: {e}") + self.client = None + raise + + async def parse_single_source(self, source: dict, fetch_count: int, min_text_length: int) -> None: + source_id = source["id"] + username = source["external_id"] + + logger.info(f"Parsing Instagram source: {username}") + + try: + user_id = await asyncio.to_thread(self.client.user_id_from_username, username) + medias = await asyncio.to_thread(self.client.user_medias, user_id, fetch_count) + except (UserNotFound, PrivateAccount) as e: + logger.warning(f"Source {username} is inaccessible: {e}") + await self.deactivate_source(source_id, str(e)) + return + except Exception as e: + logger.error(f"Error fetching {username}: {e}") + raise + + if not medias: + await self.mark_source_ok(source_id) + return + + candidates = [m.model_dump() for m in medias] + pks = [str(m["pk"]) for m in candidates] + + known_pks = await self.get_known_post_ids(source_id, pks) + candidates = [c for c in candidates if str(c["pk"]) not in known_pks] + + candidate_hashes = [make_content_hash(c.get("caption_text") or "") for c in candidates] + known_hashes = await self.get_known_hashes(candidate_hashes) + + saved = 0 + max_seen_dt = None + + for post in candidates: + try: + post_dt = post.get("taken_at") + if post_dt: + if post_dt.tzinfo is None: + post_dt = post_dt.replace(tzinfo=timezone.utc) + if max_seen_dt is None or post_dt > max_seen_dt: + max_seen_dt = post_dt + + text = (post.get("caption_text") or "").strip() + content_hash = make_content_hash(text) + + if content_hash in known_hashes: + continue + + media_type = post.get("media_type") + has_media = media_type in (1, 2, 8) + + if not has_media: + if await self.save_post_and_media(source_id, post, POST_STATUS_TRASH, "no_media"): + saved += 1 + continue + + if len(text) < min_text_length: + if await self.save_post_and_media(source_id, post, POST_STATUS_TRASH, "text_too_short"): + saved += 1 + continue + + if await self.save_post_and_media(source_id, post, POST_STATUS_MEDIA_PENDING, None): + saved += 1 + except Exception as e: + logger.error(f"Failed to save insta post {post.get('pk')}: {e}") + + await self.mark_source_ok(source_id, max_seen_dt) + logger.info(f"Source {username}: saved {saved} posts") + + async def run_once(self): + login = str(await fetch_setting("insta_login", "")).strip() + password = str(await fetch_setting("insta_password", "")).strip() + proxy = str(await fetch_setting("insta_proxy_url", "")).strip() or None + fetch_count = await fetch_int_setting("insta_fetch_count", 5) + delay_base = await fetch_int_setting("insta_delay_base_minutes", 35) + delay_random = await fetch_int_setting("insta_delay_random_minutes", 5) + cooldown_hours = await fetch_int_setting("insta_cooldown_hours", 12) + min_text_length = await fetch_int_setting("min_text_length", settings.min_text_length) + + if not login or not password: + logger.warning("Instagram login/password not set in app_settings. Skipping.") + return + + cooldown_until = await self.get_cooldown_until() + if cooldown_until and cooldown_until > datetime.now(tz=timezone.utc): + logger.info(f"Insta parser is in cooldown until {cooldown_until}. Skipping.") + return + + try: + await asyncio.to_thread(self.setup_client, login, password, proxy) + except ChallengeRequired as e: + error_msg = f"Challenge Required! Cannot login. Need manual verification.\n{e}" + logger.error(error_msg) + await send_parser_error_alert("Instagram Parser", "Login Process", error_msg) + await self.save_cooldown(cooldown_hours) + return + except Exception as e: + logger.error(f"Login setup failed: {e}") + return + + source = await self.get_active_source() + if not source: + logger.debug("No active instagram sources found.") + return + + try: + await self.parse_single_source(source, fetch_count, min_text_length) + + delay = delay_base * 60 + random.uniform(-delay_random * 60, delay_random * 60) + delay = max(60, delay) + logger.info(f"Sleeping for {delay/60:.1f} minutes before next parse.") + await asyncio.sleep(delay) + + except ChallengeRequired as e: + error_msg = f"Challenge Required during parsing!\n{e}" + logger.error(error_msg) + await send_parser_error_alert("Instagram Parser", source.get("external_id", "Unknown"), error_msg) + await self.save_cooldown(cooldown_hours) + except PleaseWaitFewMinutes as e: + error_msg = f"Rate limit hit! Cooling down.\n{e}" + logger.warning(error_msg) + await send_parser_error_alert("Instagram Parser", source.get("external_id", "Unknown"), error_msg) + await self.save_cooldown(cooldown_hours) + except Exception as e: + logger.error(f"Unexpected error parsing {source.get('external_id')}: {e}") + await asyncio.sleep(60) + + async def _heartbeat_loop(self) -> None: + while True: + try: + await self.heartbeat.beat(self.pool, {"state": "polling"}) + except Exception as e: + logger.warning(f"Insta-parser heartbeat failed: {e}") + await asyncio.sleep(self.heartbeat_interval_sec) + + async def run_loop(self): + await self.init() + logger.info("Instagram Parser started") + await self.heartbeat.beat(self.pool, {"state": "started"}, force=True) + self.heartbeat_task = asyncio.create_task(self._heartbeat_loop()) + + try: + while True: + try: + await self.heartbeat.beat(self.pool, {"state": "running"}) + await self.run_once() + except Exception as e: + logger.exception(f"Insta-parser loop error: {e}") + + await asyncio.sleep(10) + finally: + if self.heartbeat_task and not self.heartbeat_task.done(): + self.heartbeat_task.cancel() + try: + await self.heartbeat_task + except asyncio.CancelledError: + pass + + +async def main(): + worker = InstaParserWorker() + await worker.run_loop() + + +if __name__ == "__main__": + import sys + logger.remove() + logger.add( + sys.stdout, + level=settings.log_level, + format=( + "{time:YYYY-MM-DD HH:mm:ss} | " + "{level: <8} | " + "{name}:{line} — {message}" + ), + ) + asyncio.run(main())