Initial commit for RedAirsoft VK to TG and MAX poster
This commit is contained in:
+306
@@ -0,0 +1,306 @@
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import os
|
||||
import signal
|
||||
import sys
|
||||
from typing import Any, Optional
|
||||
from loguru import logger
|
||||
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.")
|
||||
Reference in New Issue
Block a user