from __future__ import annotations import asyncio import logging from aiogram import Bot, Dispatcher, Router from aiogram.types import CallbackQuery, Message from aiogram.client.session.aiohttp import AiohttpSession from apscheduler.schedulers.asyncio import AsyncIOScheduler from app.ai import OfferReplyGenerator from app.config import load_settings from app.service import OfferProcessingService from app.sources import KworkSource # from app.sources import FlSource from app.state import PublishedOffersStore, SentOffersStore from app.telegram import TelegramPublisher logging.basicConfig( level=logging.DEBUG, format="%(asctime)s | %(levelname)s | %(name)s | %(message)s", ) logger = logging.getLogger(__name__) def _build_direct_message_text(message: Message, offer_url: str, reply_text: str) -> str: message_link = _build_message_link(message) if message_link: return f"Пост: {message_link}\nОффер: {offer_url}\n\n{reply_text}" return f"Оффер: {offer_url}\n\n{reply_text}" def _build_message_link(message: Message) -> str | None: if message.chat.username: return f"https://t.me/{message.chat.username}/{message.message_id}" chat_id = str(message.chat.id) if not chat_id.startswith("-100"): return None internal_chat_id = chat_id[4:] if not internal_chat_id: return None return f"https://t.me/c/{internal_chat_id}/{message.message_id}" def build_dispatcher( reply_generator: OfferReplyGenerator, published_offers_store: PublishedOffersStore, ) -> Dispatcher: router = Router() @router.callback_query(lambda callback: callback.data == "accept_offer") async def accept_offer_callback(callback: CallbackQuery, bot: Bot) -> None: if callback.message is None: await callback.answer() return offer = published_offers_store.get(callback.message.chat.id, callback.message.message_id) if offer is None: await callback.answer("Не удалось найти оффер", show_alert=True) return await callback.answer("Готовлю сообщение") try: reply_text = await reply_generator.generate_reply(offer) await bot.send_message( chat_id=callback.from_user.id, text=_build_direct_message_text(callback.message, offer.url, reply_text), disable_web_page_preview=True, ) logger.info( "Generated DM for user_id=%s from chat_id=%s message_id=%s", callback.from_user.id, callback.message.chat.id, callback.message.message_id, ) except Exception: logger.exception( "Failed to generate or send DM for user_id=%s chat_id=%s message_id=%s", callback.from_user.id, callback.message.chat.id, callback.message.message_id, ) try: await bot.send_message( chat_id=callback.from_user.id, text="Не удалось подготовить сообщение. Убедитесь, что вы начали диалог с ботом, и попробуйте снова.", ) except Exception: logger.exception("Failed to send fallback DM to user_id=%s", callback.from_user.id) @router.callback_query(lambda callback: callback.data == "delete_post") async def delete_post_callback(callback: CallbackQuery) -> None: if callback.message is None: await callback.answer() return try: await callback.message.delete() published_offers_store.remove(callback.message.chat.id, callback.message.message_id) logger.info( "Deleted Telegram message chat_id=%s message_id=%s by user_id=%s", callback.message.chat.id, callback.message.message_id, callback.from_user.id, ) await callback.answer("Пост удален") except Exception: logger.exception( "Failed to delete Telegram message chat_id=%s message_id=%s", callback.message.chat.id, callback.message.message_id, ) await callback.answer("Не удалось удалить пост", show_alert=True) dispatcher = Dispatcher() dispatcher.include_router(router) return dispatcher async def main() -> None: settings = load_settings() session = AiohttpSession(proxy=settings.telegram_proxy_url) if settings.telegram_proxy_url else None bot = Bot(token=settings.telegram_bot_token, session=session) publisher = TelegramPublisher(bot=bot, channel_id=settings.telegram_channel_id) store = SentOffersStore(settings.state_file) published_offers_store = PublishedOffersStore(settings.published_offers_file) reply_generator = OfferReplyGenerator( api_base_url=settings.openai_api_base_url, api_key=settings.openai_api_key, model=settings.openai_model, system_prompt_file=settings.openai_system_prompt_file, timeout_seconds=settings.requests_timeout_seconds, proxy_url=settings.telegram_proxy_url, ) dispatcher = build_dispatcher(reply_generator, published_offers_store) sources = [ KworkSource( category_ids=settings.kwork_category_ids, timeout_seconds=settings.requests_timeout_seconds, ) ] # if settings.fl_rss_url: # sources.append( # FlSource( # rss_url=settings.fl_rss_url, # timeout_seconds=settings.requests_timeout_seconds, # ) # ) service = OfferProcessingService( sources=sources, publisher=publisher, sent_offers_store=store, published_offers_store=published_offers_store, include_keywords=settings.keywords_include, exclude_keywords=settings.keywords_exclude, min_price_amount=settings.min_price_amount, ) scheduler = AsyncIOScheduler() scheduler.add_job(service.run_once, "interval", minutes=settings.scheduler_interval_minutes) scheduler.start() await service.run_once() try: await dispatcher.start_polling(bot) finally: scheduler.shutdown() await bot.session.close() if __name__ == "__main__": asyncio.run(main())