diff --git a/README.md b/README.md index 8f80181..5878b63 100644 --- a/README.md +++ b/README.md @@ -83,7 +83,9 @@ 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 | +| `HTTPS_PROXY` | Optional scraper proxy; can be toggled in the Web UI | empty | + +Telegram long polling uses a direct connection so bot commands arrive in real time even when the scraper proxy is enabled. ## Architecture diff --git a/worker/src/main.py b/worker/src/main.py index a52955b..bf79dd5 100644 --- a/worker/src/main.py +++ b/worker/src/main.py @@ -4,46 +4,31 @@ import logging import os import signal import sys +from collections import defaultdict from contextlib import suppress from datetime import datetime, timedelta, timezone -from typing import Any -from urllib.parse import quote import aiohttp.web as web import asyncpg +import httpx from dotenv import load_dotenv -from telegram import Update from telegram.ext import Application, ExtBot from telegram.request import HTTPXRequest from db import close_pool, get_pool from health import create_health_app, record_scheduler_run, set_start_time, set_telegram_polling -from scraper import extract_ad_fields, fetch_ads, get_network_status +from scraper import extract_ad_fields, fetch_ads from notifier import log_notification, notify_new_ad, notify_price_drop, is_user_muted, buffer_for_digest -from settings import get_turbo_mode +from settings import get_proxy_enabled, get_turbo_mode, proxy_available logger = logging.getLogger(__name__) load_dotenv() -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}" +class DirectHTTPXRequest(HTTPXRequest): + def _build_client(self) -> httpx.AsyncClient: + return httpx.AsyncClient(**self._client_kwargs, trust_env=False) def _ad_passes_filters(fields: dict, kw_row: dict) -> bool: @@ -207,12 +192,13 @@ async def flush_digests(pool: asyncpg.Pool, bot: ExtBot) -> int: async def scheduler_task(pool: object, bot: ExtBot) -> None: while True: - proxy_enabled = bool(get_network_status().get("proxy_enabled")) + proxy_enabled = False turbo_mode = False try: + proxy_enabled = proxy_available() and await get_proxy_enabled(pool) turbo_mode = await get_turbo_mode(pool) except Exception: - logger.exception("Could not load turbo mode setting") + logger.exception("Could not load scheduler mode settings") turbo_active = turbo_mode and proxy_enabled speed_divisor = 10 if turbo_active else 1 @@ -449,25 +435,23 @@ async def main() -> None: except Exception: logger.exception("Failed to start Web UI (optional)") - proxy_url = _get_proxy_url() - bot_request = HTTPXRequest( + bot_request = DirectHTTPXRequest( connection_pool_size=10, - proxy_url=proxy_url, read_timeout=30.0, write_timeout=30.0, connect_timeout=30.0, pool_timeout=10.0, media_write_timeout=60.0, ) - updates_request = HTTPXRequest( + updates_request = DirectHTTPXRequest( connection_pool_size=10, - proxy_url=proxy_url, read_timeout=30.0, write_timeout=30.0, connect_timeout=30.0, pool_timeout=10.0, media_write_timeout=60.0, ) + logger.info("Telegram Bot API traffic uses direct network path") app = ( Application.builder() diff --git a/worker/src/migrations/08-proxy-setting.sql b/worker/src/migrations/08-proxy-setting.sql new file mode 100644 index 0000000..70864e8 --- /dev/null +++ b/worker/src/migrations/08-proxy-setting.sql @@ -0,0 +1,3 @@ +INSERT INTO app_settings (key, value) +VALUES ('proxy_enabled', 'true') +ON CONFLICT (key) DO NOTHING; diff --git a/worker/src/scraper.py b/worker/src/scraper.py index 6a289e3..1bfb15a 100644 --- a/worker/src/scraper.py +++ b/worker/src/scraper.py @@ -1,38 +1,36 @@ import asyncio import logging import os -from datetime import datetime, timezone +from datetime import datetime from typing import Any -from urllib.parse import quote import httpx +from db import get_pool +from settings import get_effective_proxy_url, get_proxy_url_from_env, proxy_available + logger = logging.getLogger(__name__) _client: httpx.AsyncClient | None = None +_client_proxy_url: str | None = None _proxy_ip_logged = False _system_public_ip: str | None = None _proxy_public_ip: str | None = None _last_used_public_ip: str | None = None +_proxy_enabled_effective: bool | None = None def _get_proxy_url() -> str | None: - raw = (os.getenv("HTTPS_PROXY") or os.getenv("https_proxy") or "").strip() - if not raw: - return None + return get_proxy_url_from_env() - 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 +async def _get_effective_proxy_url() -> str | None: + global _proxy_enabled_effective - host, port, username, password = parts - return f"http://{quote(username, safe='')}:{quote(password, safe='')}@{host}:{port}" + pool = await get_pool() + proxy_url = await get_effective_proxy_url(pool) + _proxy_enabled_effective = proxy_url is not None + return proxy_url def _redact_proxy_url(proxy_url: str) -> str: @@ -88,12 +86,22 @@ async def _log_proxy_ip_comparison(proxy_url: str) -> None: async def get_client() -> httpx.AsyncClient: """Return a shared AsyncClient with keepalive connection pool.""" - global _client, _system_public_ip, _last_used_public_ip + global _client, _client_proxy_url, _system_public_ip, _last_used_public_ip + + proxy_url = await _get_effective_proxy_url() + + if _client is not None and not _client.is_closed and _client_proxy_url != proxy_url: + await _client.aclose() + logger.info( + "Recreated httpx client because proxy changed: %s -> %s", + "enabled" if _client_proxy_url else "disabled", + "enabled" if proxy_url else "disabled", + ) + _client = None 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) @@ -115,6 +123,7 @@ async def get_client() -> httpx.AsyncClient: proxy=proxy_url, trust_env=False, ) + _client_proxy_url = proxy_url logger.info( "Created httpx client: max_conns=%d, keepalive=%d, proxy=%s", max_conns, @@ -127,9 +136,10 @@ async def get_client() -> httpx.AsyncClient: async def refresh_network_status() -> dict[str, str | bool | None]: """Ensure network status has best-effort values even before first scrape.""" - global _system_public_ip, _proxy_public_ip, _last_used_public_ip + global _system_public_ip, _proxy_public_ip, _last_used_public_ip, _proxy_enabled_effective - proxy_url = _get_proxy_url() + proxy_url = await _get_effective_proxy_url() + _proxy_enabled_effective = proxy_url is not None if _system_public_ip is None: direct_client = httpx.AsyncClient(timeout=10.0, trust_env=False) @@ -155,12 +165,14 @@ async def refresh_network_status() -> dict[str, str | bool | None]: def get_network_status() -> dict[str, str | bool | None]: """Return best-effort network status for UI display.""" - proxy_enabled = bool(_get_proxy_url()) + available = proxy_available() + proxy_enabled = bool(_proxy_enabled_effective) if _proxy_enabled_effective is not None else available last_used = _last_used_public_ip if last_used is None: last_used = _proxy_public_ip if proxy_enabled else _system_public_ip return { + "proxy_available": available, "proxy_enabled": proxy_enabled, "system_public_ip": _system_public_ip, "proxy_public_ip": _proxy_public_ip, @@ -170,11 +182,12 @@ def get_network_status() -> dict[str, str | bool | None]: async def close_client() -> None: """Close the shared AsyncClient. Call during shutdown.""" - global _client + global _client, _client_proxy_url if _client and not _client.is_closed: await _client.aclose() logger.info("Closed httpx client") - _client = None + _client = None + _client_proxy_url = None _API_URL = ( diff --git a/worker/src/settings.py b/worker/src/settings.py index 6d9c78d..fb193d7 100644 --- a/worker/src/settings.py +++ b/worker/src/settings.py @@ -1,5 +1,42 @@ +import logging +import os +from urllib.parse import quote + import asyncpg +logger = logging.getLogger(__name__) + +TRUE_VALUES = {"1", "true", "yes", "on"} + + +def _as_bool(raw: object, default: bool = False) -> bool: + if raw is None: + return default + return str(raw).lower() in TRUE_VALUES + + +def get_proxy_url_from_env() -> 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 proxy_available() -> bool: + return get_proxy_url_from_env() is not None + async def ensure_app_settings(pool: asyncpg.Pool) -> None: await pool.execute( @@ -13,20 +50,43 @@ async def ensure_app_settings(pool: asyncpg.Pool) -> None: ) -async def get_turbo_mode(pool: asyncpg.Pool) -> bool: +async def get_bool_setting(pool: asyncpg.Pool, key: str, default: bool = False) -> bool: await ensure_app_settings(pool) - raw = await pool.fetchval("SELECT value FROM app_settings WHERE key = 'turbo_mode'") - return str(raw).lower() in {"1", "true", "yes", "on"} + raw = await pool.fetchval("SELECT value FROM app_settings WHERE key = $1", key) + return _as_bool(raw, default=default) -async def set_turbo_mode(pool: asyncpg.Pool, enabled: bool) -> None: +async def set_bool_setting(pool: asyncpg.Pool, key: str, enabled: bool) -> None: await ensure_app_settings(pool) await pool.execute( """ INSERT INTO app_settings (key, value, updated_at) - VALUES ('turbo_mode', $1, now()) + VALUES ($1, $2, now()) ON CONFLICT (key) DO UPDATE SET value = EXCLUDED.value, updated_at = now() """, + key, "true" if enabled else "false", ) + + +async def get_turbo_mode(pool: asyncpg.Pool) -> bool: + return await get_bool_setting(pool, "turbo_mode") + + +async def set_turbo_mode(pool: asyncpg.Pool, enabled: bool) -> None: + await set_bool_setting(pool, "turbo_mode", enabled) + + +async def get_proxy_enabled(pool: asyncpg.Pool) -> bool: + return await get_bool_setting(pool, "proxy_enabled", default=proxy_available()) + + +async def set_proxy_enabled(pool: asyncpg.Pool, enabled: bool) -> None: + await set_bool_setting(pool, "proxy_enabled", enabled) + + +async def get_effective_proxy_url(pool: asyncpg.Pool) -> str | None: + if not await get_proxy_enabled(pool): + return None + return get_proxy_url_from_env() diff --git a/worker/src/templates/dashboard.html b/worker/src/templates/dashboard.html index 13d935d..738d4f4 100644 --- a/worker/src/templates/dashboard.html +++ b/worker/src/templates/dashboard.html @@ -46,14 +46,25 @@