Files
kwork/main.py

182 lines
6.4 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
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())