317 lines
12 KiB
Python
317 lines
12 KiB
Python
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import os
|
|
import signal
|
|
import sys
|
|
from typing import Any, Optional
|
|
from loguru import logger
|
|
try:
|
|
from .admin_notifier import AdminNotifier
|
|
from .cleaner import cleanup_stale_cache, run_cleaner_loop
|
|
from .config import settings
|
|
from .database import Database
|
|
from .max_poster import MAXPoster
|
|
from .media_processor import MediaProcessor
|
|
from .tg_poster import TelegramPoster
|
|
from .vk_client import VKClient, VKPost
|
|
except (ImportError, ValueError):
|
|
from admin_notifier import AdminNotifier
|
|
from cleaner import cleanup_stale_cache, run_cleaner_loop
|
|
from config import settings
|
|
from database import Database
|
|
from max_poster import MAXPoster
|
|
from media_processor import MediaProcessor
|
|
from tg_poster import TelegramPoster
|
|
from vk_client import VKClient, VKPost
|
|
|
|
|
|
class ServiceApp:
|
|
def __init__(self) -> None:
|
|
self.db = Database()
|
|
self.tg_poster = TelegramPoster()
|
|
self.max_poster = MAXPoster()
|
|
self.admin_notifier: Optional[AdminNotifier] = None
|
|
self.vk_group_owner_id: Optional[int] = None
|
|
self.vk_group_name: str = ""
|
|
self.vk_group_url: str = ""
|
|
self.running = False
|
|
self.cleaner_task: Optional[asyncio.Task] = None
|
|
|
|
async def init(self) -> None:
|
|
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>",
|
|
)
|
|
logger.info("Initializing VK to TG & MAX Poster Service...")
|
|
|
|
# Initialize SQLite DB
|
|
await self.db.init()
|
|
|
|
# Initialize Telegram Poster
|
|
await self.tg_poster.init()
|
|
if self.tg_poster.bot:
|
|
self.admin_notifier = AdminNotifier(self.tg_poster.bot)
|
|
|
|
# Resolve VK Group
|
|
if not settings.vk_source:
|
|
raise ValueError("VK_SOURCE is not set in configuration")
|
|
|
|
async with VKClient() as vk:
|
|
screen_name, owner_id, name = await vk.resolve_group(settings.vk_source)
|
|
self.vk_group_owner_id = owner_id
|
|
self.vk_group_name = name
|
|
self.vk_group_url = f"https://vk.com/{screen_name}"
|
|
logger.info("Resolved VK Group: '{}' (owner_id: {}, url: {})", name, owner_id, self.vk_group_url)
|
|
|
|
# Start periodic background cache cleaner
|
|
self.cleaner_task = asyncio.create_task(run_cleaner_loop(interval_minutes=15))
|
|
|
|
async def process_new_post(self, post: VKPost) -> dict[str, Any]:
|
|
vk_url = f"https://vk.com/wall{post.owner_id}_{post.post_id}"
|
|
logger.info("Processing post #{} from {}", post.post_id, vk_url)
|
|
|
|
# Save to database
|
|
post_db_id = await self.db.save_or_update_post(
|
|
owner_id=post.owner_id,
|
|
post_id=post.post_id,
|
|
posted_at=post.date,
|
|
text=post.text,
|
|
raw_data=post.raw,
|
|
)
|
|
|
|
media_processor = MediaProcessor(is_local_tg_api=self.tg_poster.is_local_api)
|
|
processed_media = []
|
|
|
|
result_summary: dict[str, Any] = {
|
|
"vk_post_id": post.post_id,
|
|
"vk_post_url": vk_url,
|
|
"tg_status": "pending",
|
|
"tg_url": None,
|
|
"tg_error": None,
|
|
"max_status": "pending",
|
|
"max_url": None,
|
|
"max_error": None,
|
|
}
|
|
|
|
try:
|
|
# 1. Download/extract media
|
|
if post.media:
|
|
logger.info("Downloading {} media items for post #{}...", len(post.media), post.post_id)
|
|
processed_media = await media_processor.process_media_items(post.media)
|
|
|
|
# 2. Publish to Telegram
|
|
try:
|
|
tg_mids, tg_url = await self.tg_poster.post_to_telegram(
|
|
raw_text=post.text,
|
|
media_items=processed_media,
|
|
vk_url=vk_url,
|
|
)
|
|
await self.db.update_tg_result(
|
|
post_db_id=post_db_id,
|
|
status="published",
|
|
message_ids=tg_mids,
|
|
url=tg_url,
|
|
)
|
|
result_summary["tg_status"] = "published"
|
|
result_summary["tg_url"] = tg_url
|
|
except Exception as exc:
|
|
err = str(exc)
|
|
logger.exception("Telegram post error for #{}: {}", post.post_id, exc)
|
|
await self.db.update_tg_result(post_db_id=post_db_id, status="failed", error=err)
|
|
result_summary["tg_status"] = "failed"
|
|
result_summary["tg_error"] = err
|
|
|
|
# 3. Publish to MAX Messenger
|
|
if settings.max_bot_token and settings.max_chat_id:
|
|
try:
|
|
max_mids, max_url = await self.max_poster.post_to_max(
|
|
raw_text=post.text,
|
|
media_items=processed_media,
|
|
vk_url=vk_url,
|
|
)
|
|
await self.db.update_max_result(
|
|
post_db_id=post_db_id,
|
|
status="published",
|
|
message_ids=max_mids,
|
|
url=max_url,
|
|
)
|
|
result_summary["max_status"] = "published"
|
|
result_summary["max_url"] = max_url
|
|
except Exception as exc:
|
|
err = str(exc)
|
|
logger.exception("MAX post error for #{}: {}", post.post_id, exc)
|
|
await self.db.update_max_result(post_db_id=post_db_id, status="failed", error=err)
|
|
result_summary["max_status"] = "failed"
|
|
result_summary["max_error"] = err
|
|
else:
|
|
await self.db.update_max_result(post_db_id=post_db_id, status="skipped")
|
|
result_summary["max_status"] = "skipped"
|
|
|
|
finally:
|
|
# Immediate cleanup of temporary media files
|
|
if processed_media:
|
|
await media_processor.cleanup(processed_media)
|
|
|
|
return result_summary
|
|
|
|
async def run_cycle(self) -> None:
|
|
if self.vk_group_owner_id is None:
|
|
return
|
|
|
|
logger.info("Checking VK group '{}' for new posts...", self.vk_group_name)
|
|
posts_to_process: list[VKPost] = []
|
|
cycle_error: Optional[str] = None
|
|
published_reports: list[dict[str, Any]] = []
|
|
|
|
try:
|
|
async with VKClient() as vk:
|
|
latest_posts = await vk.get_latest_posts(
|
|
owner_id=self.vk_group_owner_id,
|
|
count=settings.vk_check_count,
|
|
)
|
|
|
|
# Check if this is the very first run on an empty database
|
|
is_initial_start = not await self.db.has_any_posts(self.vk_group_owner_id)
|
|
if is_initial_start and latest_posts:
|
|
logger.info(
|
|
"Initial start detected on fresh database. Applying bootstrap mode: '{}'",
|
|
settings.bootstrap_mode,
|
|
)
|
|
if settings.bootstrap_mode == "skip_existing":
|
|
for p in latest_posts:
|
|
await self.db.mark_post_skipped(
|
|
owner_id=p.owner_id,
|
|
post_id=p.post_id,
|
|
posted_at=p.date,
|
|
text=p.text,
|
|
raw_data=p.raw,
|
|
reason="initial_bootstrap_skip",
|
|
)
|
|
logger.info("Marked {} existing posts as already known. Only new future posts will be published.", len(latest_posts))
|
|
latest_posts = []
|
|
elif settings.bootstrap_mode == "publish_latest_one":
|
|
# Mark all except the single latest post as skipped
|
|
newest = max(latest_posts, key=lambda p: p.date)
|
|
for p in latest_posts:
|
|
if p.post_id != newest.post_id:
|
|
await self.db.mark_post_skipped(
|
|
owner_id=p.owner_id,
|
|
post_id=p.post_id,
|
|
posted_at=p.date,
|
|
text=p.text,
|
|
raw_data=p.raw,
|
|
reason="initial_bootstrap_skip",
|
|
)
|
|
latest_posts = [newest]
|
|
logger.info("Bootstrap mode: keeping only the single newest post #{}", newest.post_id)
|
|
|
|
# Filter out processed posts and sort oldest -> newest
|
|
for p in latest_posts:
|
|
if p.is_repost:
|
|
logger.debug("Skipping repost #{}", p.post_id)
|
|
continue
|
|
is_done = await self.db.is_post_processed(p.owner_id, p.post_id)
|
|
if not is_done:
|
|
posts_to_process.append(p)
|
|
|
|
# Sort chronological (oldest to newest)
|
|
posts_to_process.sort(key=lambda p: p.date)
|
|
|
|
if posts_to_process:
|
|
logger.info("Found {} new posts to publish", len(posts_to_process))
|
|
for post in posts_to_process:
|
|
report = await self.process_new_post(post)
|
|
published_reports.append(report)
|
|
await asyncio.sleep(2.0) # Pause between posts
|
|
else:
|
|
logger.info("No new posts found.")
|
|
|
|
tg_success = sum(1 for r in published_reports if r.get("tg_status") == "published")
|
|
max_success = sum(1 for r in published_reports if r.get("max_status") == "published")
|
|
await self.db.record_run(
|
|
found_count=len(posts_to_process),
|
|
tg_count=tg_success,
|
|
max_count=max_success,
|
|
status="ok",
|
|
)
|
|
|
|
except Exception as exc:
|
|
cycle_error = str(exc)
|
|
logger.exception("Error during parse cycle: {}", exc)
|
|
await self.db.record_run(
|
|
found_count=0,
|
|
tg_count=0,
|
|
max_count=0,
|
|
status="error",
|
|
error=cycle_error,
|
|
)
|
|
|
|
# Notify admins
|
|
if self.admin_notifier:
|
|
await self.admin_notifier.notify_cycle_result(
|
|
group_name=self.vk_group_name,
|
|
vk_url=self.vk_group_url,
|
|
found_posts=published_reports,
|
|
error=cycle_error,
|
|
)
|
|
|
|
async def run(self) -> None:
|
|
await self.init()
|
|
self.running = True
|
|
logger.info("Service started. Checking every {} minutes.", settings.check_interval_minutes)
|
|
|
|
while self.running:
|
|
try:
|
|
await self.run_cycle()
|
|
except Exception as exc:
|
|
logger.exception("Unexpected error in main loop: {}", exc)
|
|
|
|
# Sleep between cycles
|
|
sleep_seconds = settings.check_interval_minutes * 60
|
|
logger.debug("Sleeping for {} seconds until next check...", sleep_seconds)
|
|
for _ in range(sleep_seconds):
|
|
if not self.running:
|
|
break
|
|
await asyncio.sleep(1.0)
|
|
|
|
async def close(self) -> None:
|
|
logger.info("Stopping service...")
|
|
self.running = False
|
|
if self.cleaner_task:
|
|
self.cleaner_task.cancel()
|
|
await self.tg_poster.close()
|
|
# Clean any remaining stale files on shutdown
|
|
await cleanup_stale_cache(max_age_minutes=0)
|
|
logger.info("Service stopped cleanly.")
|
|
|
|
|
|
async def main() -> None:
|
|
app = ServiceApp()
|
|
loop = asyncio.get_running_loop()
|
|
|
|
def handle_signal():
|
|
logger.info("Signal received, stopping...")
|
|
asyncio.create_task(app.close())
|
|
|
|
for sig in (signal.SIGINT, signal.SIGTERM):
|
|
try:
|
|
loop.add_signal_handler(sig, handle_signal)
|
|
except NotImplementedError:
|
|
# Signal handlers not implemented on Windows event loop for non-main threads
|
|
pass
|
|
|
|
try:
|
|
await app.run()
|
|
finally:
|
|
await app.close()
|
|
|
|
|
|
if __name__ == "__main__":
|
|
try:
|
|
asyncio.run(main())
|
|
except (KeyboardInterrupt, SystemExit):
|
|
logger.info("Application exited.")
|