182 lines
6.4 KiB
Python
182 lines
6.4 KiB
Python
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())
|