feat: integrate instagrapi parser with cool-down and alerts
This commit is contained in:
@@ -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;
|
||||||
@@ -10,3 +10,4 @@ Pillow==11.3.0
|
|||||||
python-multipart==0.0.20
|
python-multipart==0.0.20
|
||||||
uvicorn[standard]==0.35.0
|
uvicorn[standard]==0.35.0
|
||||||
yt-dlp==2026.6.9
|
yt-dlp==2026.6.9
|
||||||
|
instagrapi==2.1.2
|
||||||
|
|||||||
@@ -22,7 +22,7 @@ from fastapi.templating import Jinja2Templates
|
|||||||
from loguru import logger
|
from loguru import logger
|
||||||
|
|
||||||
from .config import settings
|
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 .db import fetch_int_setting, fetch_setting, get_pool
|
||||||
from .security import hash_password, new_token, token_hash, verify_password
|
from .security import hash_password, new_token, token_hash, verify_password
|
||||||
from .text_utils import build_publication_text, normalize_hash_tag, parse_categories
|
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.tg_reactor import TelegramReactor
|
||||||
from .workers.vk_poster import VKPoster
|
from .workers.vk_poster import VKPoster
|
||||||
from .workers.vk_storage_uploader import TelegramStorageUploader
|
from .workers.vk_storage_uploader import TelegramStorageUploader
|
||||||
|
from .workers.insta_parser import InstaParserWorker
|
||||||
|
|
||||||
COOKIE_NAME = "vk_parser_admin"
|
COOKIE_NAME = "vk_parser_admin"
|
||||||
VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier"
|
VK_OAUTH_VERIFIER_COOKIE = "vk_oauth_verifier"
|
||||||
@@ -1811,6 +1812,7 @@ async def startup() -> None:
|
|||||||
("tg-reactor", TelegramReactor()),
|
("tg-reactor", TelegramReactor()),
|
||||||
("vk-poster", VKPoster()),
|
("vk-poster", VKPoster()),
|
||||||
("vk-storage-uploader", TelegramStorageUploader()),
|
("vk-storage-uploader", TelegramStorageUploader()),
|
||||||
|
("insta-parser", InstaParserWorker()),
|
||||||
]
|
]
|
||||||
for name, worker in workers:
|
for name, worker in workers:
|
||||||
asyncio.create_task(start_worker_task(worker, name))
|
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"):
|
if not isinstance(item, dict) or not item.get("ok"):
|
||||||
continue
|
continue
|
||||||
try:
|
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(
|
row = await pool.fetchrow(
|
||||||
"""
|
"""
|
||||||
INSERT INTO sources(platform, name, tag, url, external_id, external_owner_id, active, priority, created_by)
|
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
|
ON CONFLICT DO NOTHING
|
||||||
RETURNING id
|
RETURNING id
|
||||||
""",
|
""",
|
||||||
PLATFORM_VK,
|
platform_val,
|
||||||
str(item.get("name") or item.get("external_id") or "").strip(),
|
str(item.get("name") or external_id_val or "").strip(),
|
||||||
normalize_hash_tag(str(item.get("tag") or item.get("external_id") or ""), str(item.get("external_id") or "source")),
|
normalize_hash_tag(str(item.get("tag") or external_id_val or ""), external_id_val or "source"),
|
||||||
str(item.get("url") or "").strip(),
|
url_val,
|
||||||
str(item.get("external_id") or "").strip(),
|
external_id_val,
|
||||||
item.get("external_owner_id"),
|
item.get("external_owner_id"),
|
||||||
bool(item.get("active", True)),
|
bool(item.get("active", True)),
|
||||||
user["id"],
|
user["id"],
|
||||||
@@ -2367,10 +2379,21 @@ async def source_create(
|
|||||||
return redirect("/login")
|
return redirect("/login")
|
||||||
require_csrf(user, csrf_token)
|
require_csrf(user, csrf_token)
|
||||||
platform = platform.strip().lower() or PLATFORM_VK
|
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
|
external_owner_id = None
|
||||||
resolved_name = name.strip()
|
resolved_name = name.strip()
|
||||||
resolved_url = url.strip()
|
|
||||||
status_value = "new"
|
status_value = "new"
|
||||||
status_msg = None
|
status_msg = None
|
||||||
if platform == PLATFORM_VK and external_id:
|
if platform == PLATFORM_VK and external_id:
|
||||||
|
|||||||
@@ -1,4 +1,5 @@
|
|||||||
PLATFORM_VK = "vk"
|
PLATFORM_VK = "vk"
|
||||||
|
PLATFORM_INSTAGRAM = "instagram"
|
||||||
|
|
||||||
SOURCE_STATUS_NEW = "new"
|
SOURCE_STATUS_NEW = "new"
|
||||||
SOURCE_STATUS_OK = "ok"
|
SOURCE_STATUS_OK = "ok"
|
||||||
@@ -38,3 +39,4 @@ WORKER_VK_POSTER = "vk-poster"
|
|||||||
WORKER_MAX_POSTER = "max-poster"
|
WORKER_MAX_POSTER = "max-poster"
|
||||||
WORKER_SITE_POSTER = "site-poster"
|
WORKER_SITE_POSTER = "site-poster"
|
||||||
WORKER_DAILY_REPORT = "daily-report"
|
WORKER_DAILY_REPORT = "daily-report"
|
||||||
|
WORKER_INSTA_PARSER = "insta-parser"
|
||||||
|
|||||||
@@ -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)
|
await asyncio.sleep(float(exc.retry_after) + 1)
|
||||||
finally:
|
finally:
|
||||||
await bot.session.close()
|
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"🚨 <b>CRITICAL PARSER ERROR</b>\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()
|
||||||
|
|
||||||
|
|||||||
@@ -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=(
|
||||||
|
"<green>{time:YYYY-MM-DD HH:mm:ss}</green> | "
|
||||||
|
"<level>{level: <8}</level> | "
|
||||||
|
"<cyan>{name}</cyan>:<cyan>{line}</cyan> — <level>{message}</level>"
|
||||||
|
),
|
||||||
|
)
|
||||||
|
asyncio.run(main())
|
||||||
Reference in New Issue
Block a user