From ec3cea94c904486dd095d08f3b7cc1e76d236961 Mon Sep 17 00:00:00 2001 From: Jose Lago Date: Sun, 12 Jul 2026 21:39:21 +0200 Subject: [PATCH] feat: add support for HTTP proxy configuration and display status in dashboard --- .env.example | 3 +- README.md | 1 + worker/src/main.py | 66 +++++++++++++---------- worker/src/scraper.py | 82 +++++++++++++++++++++++++++-- worker/src/templates/dashboard.html | 9 ++++ worker/src/web.py | 5 ++ 6 files changed, 133 insertions(+), 33 deletions(-) diff --git a/.env.example b/.env.example index c622ae6..a05bad8 100644 --- a/.env.example +++ b/.env.example @@ -1,5 +1,5 @@ # Telegram Bot Token (from @BotFather) -TELEGRAM_BOT_TOKEN=8653489932:AAHhyOD1jtimE7kg0zoVCUVd3l0YEz_YJgg +TELEGRAM_BOT_TOKEN=8653489932:AAG_Ins2_z3sNHX8ZlGI4mhyzmUhWAWCZlg # Direct PostgreSQL connection POSTGRES_HOST=192.168.178.3 @@ -11,3 +11,4 @@ POSTGRES_DB=postgres # Worker Configuration DEFAULT_INTERVAL_MINUTES=60 ADMIN_TELEGRAM_IDS=298181113 # Comma-separated Telegram user IDs with admin access +HTTPS_PROXY=datacenter-de.floxy.io:1338:IPv4D_TZhQUiP9C3-ttl-0:70EDRDQUpo9Jc0a diff --git a/README.md b/README.md index 2f989d6..8f80181 100644 --- a/README.md +++ b/README.md @@ -83,6 +83,7 @@ Edit `.env` before first startup. All values are read by the worker and database | `POSTGRES_DB` | Database name | `postgres` | | `JWT_SECRET` | PostgREST JWT signing key | auto-generated default | | `DEFAULT_INTERVAL_MINUTES`| Default scrape interval per keyword | `5` | +| `HTTPS_PROXY` | Optional proxy for outbound HTTP requests | empty | ## Architecture diff --git a/worker/src/main.py b/worker/src/main.py index b125118..feece5f 100644 --- a/worker/src/main.py +++ b/worker/src/main.py @@ -372,13 +372,6 @@ async def main() -> None: pool = await get_pool() - app = Application.builder().token(os.getenv("TELEGRAM_BOT_TOKEN")).build() - - from bot import register_handlers, setup_global_commands # noqa: E402 - - await setup_global_commands(app) - register_handlers(app) - # ── Start healthcheck HTTP server ────────────────────────────── health_app = create_health_app() runner = web.AppRunner(health_app) @@ -403,7 +396,26 @@ async def main() -> None: except Exception: logger.exception("Failed to start Web UI (optional)") - scheduler = asyncio.ensure_future(scheduler_task(pool, app.bot)) + app = Application.builder().token(os.getenv("TELEGRAM_BOT_TOKEN")).build() + scheduler = None + poll_task = None + bot_started = False + + from bot import register_handlers, setup_global_commands # noqa: E402 + + try: + await setup_global_commands(app) + register_handlers(app) + await app.initialize() + await app.start() + logger.info("Bot started with long polling") + + set_telegram_polling(True) + poll_task = asyncio.ensure_future(app.updater.start_polling()) # type: ignore[attr-defined] + scheduler = asyncio.ensure_future(scheduler_task(pool, app.bot)) + bot_started = True + except Exception: + logger.exception("Telegram bot startup failed — continuing in UI-only mode") loop = asyncio.get_running_loop() stop = loop.create_future() @@ -416,36 +428,32 @@ async def main() -> None: loop.add_signal_handler(sig, _signal_handler) try: - await app.initialize() - await app.start() - logger.info("Bot started with long polling") - - set_telegram_polling(True) - poll_task = asyncio.ensure_future(app.updater.start_polling()) # type: ignore[attr-defined] - await stop logger.info("Signal received — initiating graceful shutdown...") finally: # ── Cancel scheduler with grace period ─────────────────────── - logger.info("Cancelling scheduler task...") - scheduler.cancel() - with suppress(asyncio.CancelledError): - try: - await asyncio.wait_for(scheduler, timeout=5.0) - except asyncio.TimeoutError: - logger.warning("Scheduler task did not finish within 5s — force cancelled.") + if scheduler is not None: + logger.info("Cancelling scheduler task...") + scheduler.cancel() + with suppress(asyncio.CancelledError): + try: + await asyncio.wait_for(scheduler, timeout=5.0) + except asyncio.TimeoutError: + logger.warning("Scheduler task did not finish within 5s — force cancelled.") # ── Stop Telegram polling ──────────────────────────────────── - set_telegram_polling(False) - logger.info("Stopping Telegram poller...") - poll_task.cancel() - with suppress(asyncio.CancelledError): - await poll_task + if bot_started and poll_task is not None: + set_telegram_polling(False) + logger.info("Stopping Telegram poller...") + poll_task.cancel() + with suppress(asyncio.CancelledError): + await poll_task # ── Shutdown application ───────────────────────────────────── - logger.info("Shutting down Telegram bot application...") - await app.shutdown() + if bot_started: + logger.info("Shutting down Telegram bot application...") + await app.shutdown() # ── Close health server ─────────────────────────────────────── logger.info("Stopping health check server...") diff --git a/worker/src/scraper.py b/worker/src/scraper.py index 2bd5c72..b1095aa 100644 --- a/worker/src/scraper.py +++ b/worker/src/scraper.py @@ -3,13 +3,81 @@ import logging import os from datetime import datetime, timezone from typing import Any -from urllib.parse import quote_plus +from urllib.parse import quote import httpx logger = logging.getLogger(__name__) _client: httpx.AsyncClient | None = None +_proxy_ip_logged = False + + +def _get_proxy_url() -> str | None: + raw = (os.getenv("HTTPS_PROXY") or os.getenv("https_proxy") or "").strip() + if not raw: + return None + + if "://" in raw: + return raw + + parts = raw.split(":", 3) + if len(parts) != 4: + logger.warning( + "Ignoring invalid HTTPS_PROXY value (expected host:port:user:pass or full URL)" + ) + return None + + host, port, username, password = parts + return f"http://{quote(username, safe='')}:{quote(password, safe='')}@{host}:{port}" + + +def _redact_proxy_url(proxy_url: str) -> str: + try: + parsed = httpx.URL(proxy_url) + host = parsed.host or "unknown-host" + port = parsed.port or 80 + user = parsed.username or "unknown-user" + return f"{host}:{port} (user={user}, credentials=set)" + except Exception: + return "" + + +async def _fetch_public_ip(client: httpx.AsyncClient) -> str | None: + try: + resp = await client.get("https://api.ipify.org?format=json") + resp.raise_for_status() + data = resp.json() + ip = data.get("ip") + return str(ip) if ip else None + except Exception as exc: + logger.warning("Could not resolve public IP: %s", exc) + return None + + +async def _log_proxy_ip_comparison(proxy_url: str) -> None: + global _proxy_ip_logged + + if _proxy_ip_logged: + return + + _proxy_ip_logged = True + + direct_client = httpx.AsyncClient(timeout=10.0, trust_env=False) + proxy_client = httpx.AsyncClient(timeout=10.0, trust_env=False, proxy=proxy_url) + + try: + direct_ip = await _fetch_public_ip(direct_client) + proxy_ip = await _fetch_public_ip(proxy_client) + logger.info( + "HTTPS proxy enabled: %s | public_ip_without_proxy=%s | public_ip_with_proxy=%s", + _redact_proxy_url(proxy_url), + direct_ip or "unknown", + proxy_ip or "unknown", + ) + finally: + await direct_client.aclose() + await proxy_client.aclose() async def get_client() -> httpx.AsyncClient: @@ -19,6 +87,10 @@ async def get_client() -> httpx.AsyncClient: if _client is None or _client.is_closed: max_conns = int(os.getenv("HTTP_MAX_CONNECTIONS", "10")) max_keepalive = int(os.getenv("HTTP_KEEPALIVE_CONNECTIONS", "5")) + proxy_url = _get_proxy_url() + + if proxy_url: + await _log_proxy_ip_comparison(proxy_url) _client = httpx.AsyncClient( timeout=float(os.getenv("HTTP_TIMEOUT_S", "30.0")), @@ -27,10 +99,14 @@ async def get_client() -> httpx.AsyncClient: max_keepalive_connections=max_keepalive, keepalive_expiry=60, ), + proxy=proxy_url, + trust_env=False, ) logger.info( - "Created httpx client: max_conns=%d, keepalive=%d", - max_conns, max_keepalive, + "Created httpx client: max_conns=%d, keepalive=%d, proxy=%s", + max_conns, + max_keepalive, + "enabled" if proxy_url else "disabled", ) return _client diff --git a/worker/src/templates/dashboard.html b/worker/src/templates/dashboard.html index bed103b..d7c8bae 100644 --- a/worker/src/templates/dashboard.html +++ b/worker/src/templates/dashboard.html @@ -38,6 +38,15 @@ {{ data.last_scheduler.strftime('%Y-%m-%d %H:%M:%S') if data.last_scheduler else 'Never' }} +
+
HTTP Proxy
+
+ {{ 'Enabled' if data.proxy_enabled else 'Disabled' }} +
+
+ {{ 'Outbound scraping uses the proxy' if data.proxy_enabled else 'Outbound scraping goes direct' }} +
+
{% endif %} {% endblock %} \ No newline at end of file diff --git a/worker/src/web.py b/worker/src/web.py index cbc701e..b454c03 100644 --- a/worker/src/web.py +++ b/worker/src/web.py @@ -74,6 +74,10 @@ def format_postcodes(postcodes: list | None) -> str: return ", ".join(str(p) for p in postcodes) +def _proxy_enabled() -> bool: + return bool((os.getenv("HTTPS_PROXY") or os.getenv("https_proxy") or "").strip()) + + @asynccontextmanager async def lifespan(app: FastAPI): logger.info("Web UI starting up") @@ -115,6 +119,7 @@ async def dashboard(request: Request): "queue_pending": queue_pending or 0, "queue_dead": queue_dead or 0, "last_scheduler": last_scheduler["scraped_at"] if last_scheduler else None, + "proxy_enabled": _proxy_enabled(), } return templates.TemplateResponse("dashboard.html", {"request": request, "error": None, "data": data})