feat: python worker (bot, scraper, notifier, scheduler)

This commit is contained in:
opencode
2026-06-16 19:12:35 +02:00
committed by Lago
parent d32ce26d9e
commit fbbdc5e54b
7 changed files with 729 additions and 0 deletions
+202
View File
@@ -0,0 +1,202 @@
import asyncio
import json
import logging
import os
import signal
import sys
from contextlib import suppress
from dotenv import load_dotenv
from telegram import Update
from telegram.ext import Application, ExtBot
from db import close_pool, get_pool
from scraper import extract_ad_fields, fetch_ads
from notifier import log_notification, mark_notified, notify_new_ad
logger = logging.getLogger(__name__)
load_dotenv()
async def scheduler_task(pool: object, bot: ExtBot) -> None:
while True:
try:
rows = await pool.fetch(
"""
SELECT sq.id, sq.keyword, sq.interval_minutes, u.telegram_id
FROM search_queries sq
JOIN users u ON sq.user_id = u.id
WHERE sq.is_active = true
AND (sq.last_scraped_at IS NULL OR
sq.last_scraped_at < now() - (sq.interval_minutes || ' minutes')::interval)
"""
)
for row in rows:
query_id = str(row["id"])
keyword = row["keyword"]
telegram_id = row["telegram_id"]
logger.info("Scraping keyword '%s' for query %s", keyword, query_id)
try:
ads_raw, total_hits = await fetch_ads(keyword)
new_count = 0
for ad_data in ads_raw:
fields = extract_ad_fields(ad_data)
wh_ad_id = fields["wh_ad_id"]
existing = await pool.fetchrow(
"SELECT id FROM ads WHERE wh_ad_id = $1",
wh_ad_id,
)
if not existing:
ad_row = await pool.fetchrow(
"""
INSERT INTO ads (wh_ad_id, raw_json, title, price, location, url, published_at)
VALUES ($1, $2, $3, $4, $5, $6, $7)
RETURNING id
""",
wh_ad_id,
json.dumps(ad_data),
fields["title"],
fields["price"],
fields["location"],
fields["url"],
fields.get("published_at"),
)
ad_uuid = str(ad_row["id"])
else:
ad_uuid = str(existing["id"])
existing_qa = await pool.fetchrow(
"SELECT 1 FROM query_ads WHERE search_query_id = $1 AND ad_id = $2",
query_id,
ad_uuid,
)
if not existing_qa:
await pool.execute(
"INSERT INTO query_ads (search_query_id, ad_id) VALUES ($1, $2)",
query_id,
ad_uuid,
)
user_row = await pool.fetchrow(
"SELECT id FROM users WHERE telegram_id = $1",
telegram_id,
)
user_id = str(user_row["id"]) if user_row else None
notify_fields = {**fields, "keyword": keyword}
await notify_new_ad(bot, telegram_id, notify_fields)
if user_id:
await mark_notified(pool, query_id, ad_uuid)
try:
msg_id = 0
await log_notification(pool, user_id, ad_uuid, msg_id)
except Exception:
logger.exception("Failed to log notification")
new_count += 1
logger.info(
"New ad %s found for query %s (keyword=%s)",
wh_ad_id,
query_id,
keyword,
)
await pool.execute(
"UPDATE search_queries SET last_scraped_at = now() WHERE id = $1",
query_id,
)
await pool.execute(
"""
INSERT INTO scrape_logs (search_query_id, status, ads_found, new_ads)
VALUES ($1, 'success', $2, $3)
""",
query_id,
len(ads_raw),
new_count,
)
except Exception:
logger.exception("Error scraping keyword '%s' (query %s)", keyword, query_id)
await pool.execute(
"""
INSERT INTO scrape_logs (search_query_id, status, error_message)
VALUES ($1, 'error', $2)
""",
query_id,
str(sys.exc_info()[1]),
)
await asyncio.sleep(5)
except Exception:
logger.exception("Scheduler iteration error")
await asyncio.sleep(30)
async def main() -> None:
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s",
)
if not os.getenv("TELEGRAM_BOT_TOKEN"):
logger.error("TELEGRAM_BOT_TOKEN is required")
sys.exit(1)
pool = await get_pool()
app = Application.builder().token(os.getenv("TELEGRAM_BOT_TOKEN")).build()
from bot import register_handlers # noqa: E402
register_handlers(app.bot)
scheduler = asyncio.ensure_future(scheduler_task(pool, app.bot))
loop = asyncio.get_running_loop()
stop = loop.create_future()
def _signal_handler() -> None:
if not stop.done():
stop.set_result(True)
for sig in (signal.SIGINT, signal.SIGTERM):
loop.add_signal_handler(sig, _signal_handler)
try:
await app.initialize()
await app.start()
logger.info("Bot started with long polling")
poll_task = asyncio.ensure_future(app.updater.start_polling()) # type: ignore[attr-defined]
await stop
logger.info("Shutting down...")
finally:
scheduler.cancel()
with suppress(asyncio.CancelledError):
await scheduler
poll_task.cancel()
with suppress(asyncio.CancelledError):
await poll_task
await app.shutdown()
await close_pool()
logger.info("Shutdown complete")
if __name__ == "__main__":
asyncio.run(main())