Add VPN admin tab with subscription proxy via cp-vpn for selected sources.
Persist settings in CA, expose /admin/vpn and /internal/vpn, and route Telegram (and optional web/nlp) traffic through mihomo SOCKS when enabled. Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
@@ -0,0 +1,31 @@
|
||||
FROM python:3.12-alpine
|
||||
|
||||
ARG MIHOMO_VERSION=1.19.3
|
||||
ARG TARGETARCH
|
||||
|
||||
RUN apk add --no-cache ca-certificates curl gzip \
|
||||
&& arch="${TARGETARCH:-amd64}" \
|
||||
&& case "$arch" in \
|
||||
amd64|x86_64) mihomo_arch=amd64 ;; \
|
||||
arm64|aarch64) mihomo_arch=arm64 ;; \
|
||||
*) mihomo_arch=amd64 ;; \
|
||||
esac \
|
||||
&& curl -fsSL \
|
||||
"https://github.com/MetaCubeX/mihomo/releases/download/v${MIHOMO_VERSION}/mihomo-linux-${mihomo_arch}-compatible-v${MIHOMO_VERSION}.gz" \
|
||||
-o /tmp/mihomo.gz \
|
||||
&& gunzip -c /tmp/mihomo.gz > /usr/local/bin/mihomo \
|
||||
&& chmod +x /usr/local/bin/mihomo \
|
||||
&& rm -f /tmp/mihomo.gz \
|
||||
&& mkdir -p /etc/mihomo/providers
|
||||
|
||||
WORKDIR /app
|
||||
COPY controller.py /app/controller.py
|
||||
|
||||
ENV MIHOMO_BIN=/usr/local/bin/mihomo \
|
||||
MIHOMO_CONFIG_DIR=/etc/mihomo \
|
||||
VPN_SOCKS_PORT=1080 \
|
||||
VPN_POLL_SECONDS=15
|
||||
|
||||
EXPOSE 1080
|
||||
|
||||
CMD ["python", "-u", "/app/controller.py"]
|
||||
@@ -0,0 +1,204 @@
|
||||
"""Poll CA /internal/vpn and keep mihomo config in sync for subscription mode."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import logging
|
||||
import os
|
||||
import signal
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
from pathlib import Path
|
||||
from urllib.error import HTTPError, URLError
|
||||
from urllib.request import Request, urlopen
|
||||
|
||||
logging.basicConfig(
|
||||
level=logging.INFO,
|
||||
format="%(asctime)s %(levelname)s %(message)s",
|
||||
)
|
||||
logger = logging.getLogger("cp-vpn")
|
||||
|
||||
CA_API_URL = os.getenv("CA_API_URL", "http://ca-api:8000").rstrip("/")
|
||||
INTERNAL_TOKEN = os.getenv("INTERNAL_TOKEN", "dev-internal-token")
|
||||
POLL_SECONDS = int(os.getenv("VPN_POLL_SECONDS", "15"))
|
||||
CONFIG_DIR = Path(os.getenv("MIHOMO_CONFIG_DIR", "/etc/mihomo"))
|
||||
CONFIG_PATH = CONFIG_DIR / "config.yaml"
|
||||
MIHOMO_BIN = os.getenv("MIHOMO_BIN", "/usr/local/bin/mihomo")
|
||||
MIXED_PORT = int(os.getenv("VPN_SOCKS_PORT", "1080"))
|
||||
CONTROLLER = "127.0.0.1:9090"
|
||||
|
||||
_mihomo_proc: subprocess.Popen | None = None
|
||||
_last_fingerprint: str | None = None
|
||||
|
||||
|
||||
def fetch_vpn() -> dict:
|
||||
req = Request(
|
||||
f"{CA_API_URL}/internal/vpn",
|
||||
headers={"X-Internal-Token": INTERNAL_TOKEN},
|
||||
method="GET",
|
||||
)
|
||||
with urlopen(req, timeout=15) as resp:
|
||||
return json.loads(resp.read().decode("utf-8"))
|
||||
|
||||
|
||||
def fingerprint(cfg: dict) -> str:
|
||||
return json.dumps(
|
||||
{
|
||||
"enabled": cfg.get("enabled"),
|
||||
"mode": cfg.get("mode"),
|
||||
"subscription_url": cfg.get("subscription_url"),
|
||||
"subscription_interval_seconds": cfg.get("subscription_interval_seconds"),
|
||||
},
|
||||
sort_keys=True,
|
||||
)
|
||||
|
||||
|
||||
def write_direct_config() -> None:
|
||||
CONFIG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
(CONFIG_DIR / "providers").mkdir(parents=True, exist_ok=True)
|
||||
CONFIG_PATH.write_text(
|
||||
f"""# managed by cp-vpn controller — DIRECT (subscription inactive)
|
||||
mixed-port: {MIXED_PORT}
|
||||
allow-lan: true
|
||||
bind-address: "*"
|
||||
mode: direct
|
||||
log-level: warning
|
||||
external-controller: {CONTROLLER}
|
||||
""",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
|
||||
def write_subscription_config(url: str, interval: int) -> None:
|
||||
CONFIG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
(CONFIG_DIR / "providers").mkdir(parents=True, exist_ok=True)
|
||||
# YAML: quote URL; escape double quotes in URL if any
|
||||
safe_url = url.replace("\\", "\\\\").replace('"', '\\"')
|
||||
CONFIG_PATH.write_text(
|
||||
f"""# managed by cp-vpn controller — subscription
|
||||
mixed-port: {MIXED_PORT}
|
||||
allow-lan: true
|
||||
bind-address: "*"
|
||||
mode: rule
|
||||
log-level: info
|
||||
external-controller: {CONTROLLER}
|
||||
proxy-providers:
|
||||
sub:
|
||||
type: http
|
||||
url: "{safe_url}"
|
||||
interval: {interval}
|
||||
path: ./providers/sub.yaml
|
||||
health-check:
|
||||
enable: true
|
||||
url: https://www.gstatic.com/generate_204
|
||||
interval: 600
|
||||
proxy-groups:
|
||||
- name: PROXY
|
||||
type: select
|
||||
use:
|
||||
- sub
|
||||
rules:
|
||||
- MATCH,PROXY
|
||||
""",
|
||||
encoding="utf-8",
|
||||
)
|
||||
|
||||
|
||||
def start_mihomo() -> None:
|
||||
global _mihomo_proc
|
||||
if _mihomo_proc is not None and _mihomo_proc.poll() is None:
|
||||
return
|
||||
logger.info("Starting mihomo (%s)", MIHOMO_BIN)
|
||||
_mihomo_proc = subprocess.Popen(
|
||||
[MIHOMO_BIN, "-d", str(CONFIG_DIR)],
|
||||
stdout=sys.stdout,
|
||||
stderr=sys.stderr,
|
||||
)
|
||||
|
||||
|
||||
def stop_mihomo() -> None:
|
||||
global _mihomo_proc
|
||||
if _mihomo_proc is None:
|
||||
return
|
||||
if _mihomo_proc.poll() is None:
|
||||
logger.info("Stopping mihomo")
|
||||
_mihomo_proc.send_signal(signal.SIGTERM)
|
||||
try:
|
||||
_mihomo_proc.wait(timeout=10)
|
||||
except subprocess.TimeoutExpired:
|
||||
_mihomo_proc.kill()
|
||||
_mihomo_proc = None
|
||||
|
||||
|
||||
def reload_mihomo() -> None:
|
||||
"""Force reload via external controller; fall back to process restart."""
|
||||
body = json.dumps({"path": str(CONFIG_PATH)}).encode("utf-8")
|
||||
req = Request(
|
||||
f"http://{CONTROLLER}/configs?force=true",
|
||||
data=body,
|
||||
headers={"Content-Type": "application/json"},
|
||||
method="PUT",
|
||||
)
|
||||
try:
|
||||
with urlopen(req, timeout=10) as resp:
|
||||
resp.read()
|
||||
logger.info("Mihomo config reloaded")
|
||||
return
|
||||
except (HTTPError, URLError, OSError) as exc:
|
||||
logger.warning("Reload via API failed (%s); restarting mihomo", exc)
|
||||
stop_mihomo()
|
||||
start_mihomo()
|
||||
|
||||
|
||||
def apply_config(cfg: dict) -> None:
|
||||
global _last_fingerprint
|
||||
fp = fingerprint(cfg)
|
||||
if fp == _last_fingerprint and _mihomo_proc is not None and _mihomo_proc.poll() is None:
|
||||
return
|
||||
|
||||
enabled = bool(cfg.get("enabled"))
|
||||
mode = (cfg.get("mode") or "").strip()
|
||||
url = (cfg.get("subscription_url") or "").strip()
|
||||
interval = int(cfg.get("subscription_interval_seconds") or 3600)
|
||||
|
||||
if enabled and mode == "subscription" and url:
|
||||
logger.info("Applying subscription config (interval=%ss)", interval)
|
||||
write_subscription_config(url, interval)
|
||||
else:
|
||||
logger.info("Applying DIRECT config (enabled=%s mode=%s)", enabled, mode)
|
||||
write_direct_config()
|
||||
|
||||
already_running = _mihomo_proc is not None and _mihomo_proc.poll() is None
|
||||
start_mihomo()
|
||||
if already_running or _last_fingerprint is not None:
|
||||
time.sleep(1)
|
||||
reload_mihomo()
|
||||
_last_fingerprint = fp
|
||||
|
||||
|
||||
def main() -> None:
|
||||
global _mihomo_proc
|
||||
logger.info(
|
||||
"cp-vpn controller started (poll=%ss api=%s)",
|
||||
POLL_SECONDS,
|
||||
CA_API_URL,
|
||||
)
|
||||
write_direct_config()
|
||||
start_mihomo()
|
||||
|
||||
while True:
|
||||
try:
|
||||
cfg = fetch_vpn()
|
||||
apply_config(cfg)
|
||||
except Exception:
|
||||
logger.exception("Failed to sync VPN config")
|
||||
if _mihomo_proc is not None and _mihomo_proc.poll() is not None:
|
||||
logger.warning("mihomo exited with code %s; restarting", _mihomo_proc.returncode)
|
||||
_mihomo_proc = None
|
||||
start_mihomo()
|
||||
time.sleep(POLL_SECONDS)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
@@ -179,10 +179,31 @@ async def run_with_listener() -> None:
|
||||
await worker_loop(ctx=WorkerContext())
|
||||
return
|
||||
|
||||
from workers.sources.telegram_client import TelegramAuthError, TelegramConfigError
|
||||
from workers.sources.telegram_listener import TelegramListener
|
||||
from workers.sources.telegram_session import close_shared_client, get_shared_client
|
||||
|
||||
client = await get_shared_client()
|
||||
while True:
|
||||
try:
|
||||
client = await get_shared_client()
|
||||
break
|
||||
except (TelegramAuthError, TelegramConfigError, ConnectionError, OSError) as exc:
|
||||
logger.error(
|
||||
"Telegram connect failed (%s); retry in 60s (batch without TG)",
|
||||
exc,
|
||||
)
|
||||
try:
|
||||
await asyncio.wait_for(
|
||||
worker_loop(ctx=WorkerContext(tg_client=None)),
|
||||
timeout=60,
|
||||
)
|
||||
except asyncio.TimeoutError:
|
||||
pass
|
||||
continue
|
||||
except Exception:
|
||||
logger.exception("Unexpected Telegram connect error; retry in 60s")
|
||||
await asyncio.sleep(60)
|
||||
|
||||
listener = TelegramListener(client)
|
||||
ctx = WorkerContext(tg_client=client)
|
||||
worker_task = asyncio.create_task(worker_loop(ctx=ctx))
|
||||
|
||||
@@ -57,7 +57,25 @@ class Crawl4AIAdapter:
|
||||
events: list[dict] = []
|
||||
errors: list[str] = []
|
||||
|
||||
async with AsyncWebCrawler(verbose=False) as crawler:
|
||||
crawler_kwargs: dict[str, Any] = {"verbose": False}
|
||||
try:
|
||||
from workers.sources.vpn_config import http_proxy_url
|
||||
|
||||
proxy = http_proxy_url("crawl4ai")
|
||||
if proxy:
|
||||
try:
|
||||
from crawl4ai import BrowserConfig # type: ignore
|
||||
|
||||
crawler_kwargs["config"] = BrowserConfig(proxy=proxy)
|
||||
logger.info("Crawl4AI via proxy %s", proxy.split("@")[-1])
|
||||
except Exception:
|
||||
logger.exception(
|
||||
"Could not apply VPN proxy to Crawl4AI; continuing without"
|
||||
)
|
||||
except Exception:
|
||||
logger.exception("VPN config lookup failed for crawl4ai")
|
||||
|
||||
async with AsyncWebCrawler(**crawler_kwargs) as crawler:
|
||||
for url in cfg.urls:
|
||||
try:
|
||||
if cfg.extract_mode == "llm":
|
||||
|
||||
@@ -70,7 +70,14 @@ class ViinaAdapter:
|
||||
|
||||
|
||||
async def _fetch_article_text(url: str) -> str:
|
||||
async with httpx.AsyncClient(timeout=60.0, follow_redirects=True) as client:
|
||||
from workers.sources.vpn_config import http_proxy_url
|
||||
|
||||
proxy = http_proxy_url("viina")
|
||||
kwargs: dict = {"timeout": 60.0, "follow_redirects": True}
|
||||
if proxy:
|
||||
kwargs["proxy"] = proxy
|
||||
logger.info("VIINA fetch via proxy %s", proxy.split("@")[-1])
|
||||
async with httpx.AsyncClient(**kwargs) as client:
|
||||
response = await client.get(
|
||||
url,
|
||||
headers={"User-Agent": "MapMil-CP-Viina/1.0"},
|
||||
|
||||
@@ -1,20 +1,39 @@
|
||||
"""Единое подключение Telethon для listener и batch-заданий."""
|
||||
|
||||
import logging
|
||||
from pathlib import Path
|
||||
|
||||
from telethon import TelegramClient
|
||||
from telethon.errors import AuthKeyUnregisteredError, SessionPasswordNeededError
|
||||
|
||||
from workers.sources.telegram_settings import create_client, get_api_credentials, get_session_path
|
||||
from workers.sources.telegram_settings import (
|
||||
create_client,
|
||||
describe_connection,
|
||||
get_api_credentials,
|
||||
get_session_path,
|
||||
)
|
||||
from workers.sources.telegram_client import TelegramAuthError, TelegramConfigError
|
||||
from workers.sources.vpn_config import proxy_fingerprint
|
||||
|
||||
logger = logging.getLogger("cp-worker.telegram-session")
|
||||
|
||||
_shared_client: TelegramClient | None = None
|
||||
_proxy_fingerprint: str | None = None
|
||||
|
||||
|
||||
async def get_shared_client() -> TelegramClient:
|
||||
global _shared_client
|
||||
global _shared_client, _proxy_fingerprint
|
||||
|
||||
fp = proxy_fingerprint()
|
||||
if _shared_client is not None and _shared_client.is_connected():
|
||||
return _shared_client
|
||||
if fp == _proxy_fingerprint:
|
||||
return _shared_client
|
||||
logger.info(
|
||||
"VPN proxy changed (%s → %s); reconnecting Telethon",
|
||||
_proxy_fingerprint,
|
||||
fp,
|
||||
)
|
||||
await close_shared_client()
|
||||
|
||||
try:
|
||||
api_id, api_hash = get_api_credentials()
|
||||
@@ -29,6 +48,7 @@ async def get_shared_client() -> TelegramClient:
|
||||
)
|
||||
|
||||
client = create_client(session_path, api_id, api_hash)
|
||||
logger.info("Connecting Telethon %s", describe_connection())
|
||||
try:
|
||||
await client.connect()
|
||||
if not await client.is_user_authorized():
|
||||
@@ -45,11 +65,13 @@ async def get_shared_client() -> TelegramClient:
|
||||
) from exc
|
||||
|
||||
_shared_client = client
|
||||
_proxy_fingerprint = fp
|
||||
return client
|
||||
|
||||
|
||||
async def close_shared_client() -> None:
|
||||
global _shared_client
|
||||
global _shared_client, _proxy_fingerprint
|
||||
if _shared_client is not None:
|
||||
await _shared_client.disconnect()
|
||||
_shared_client = None
|
||||
_proxy_fingerprint = None
|
||||
|
||||
@@ -28,7 +28,8 @@ def get_api_credentials() -> tuple[int, str]:
|
||||
return api_id, api_hash
|
||||
|
||||
|
||||
def get_proxy() -> tuple | None:
|
||||
def get_env_proxy() -> tuple | None:
|
||||
"""Legacy TELEGRAM_PROXY_* env fallback (used when CA VPN is off)."""
|
||||
proxy_type = os.environ.get("TELEGRAM_PROXY_TYPE", "").strip().lower()
|
||||
if not proxy_type or proxy_type == "none":
|
||||
return None
|
||||
@@ -49,6 +50,15 @@ def get_proxy() -> tuple | None:
|
||||
return proxy_type, host, port
|
||||
|
||||
|
||||
def get_proxy() -> tuple | None:
|
||||
try:
|
||||
from workers.sources.vpn_config import telethon_proxy_tuple
|
||||
|
||||
return telethon_proxy_tuple()
|
||||
except Exception:
|
||||
return get_env_proxy()
|
||||
|
||||
|
||||
def describe_connection() -> str:
|
||||
proxy = get_proxy()
|
||||
if not proxy:
|
||||
|
||||
@@ -0,0 +1,124 @@
|
||||
"""Fetch VPN settings from CA /internal/vpn (cached)."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
import os
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
|
||||
from contracts.vpn import VpnSettings
|
||||
|
||||
logger = logging.getLogger("cp-worker.vpn")
|
||||
|
||||
CA_API_URL = os.getenv("CA_API_URL", "http://ca-api:8000").rstrip("/")
|
||||
INTERNAL_TOKEN = os.getenv("INTERNAL_TOKEN", "dev-internal-token")
|
||||
CACHE_TTL = float(os.getenv("VPN_CONFIG_CACHE_SECONDS", "30"))
|
||||
CP_VPN_HOST = os.getenv("CP_VPN_HOST", "cp-vpn")
|
||||
CP_VPN_PORT = int(os.getenv("CP_VPN_PORT", "1080"))
|
||||
|
||||
_cache: VpnSettings | None = None
|
||||
_cache_at: float = 0.0
|
||||
_cache_failed: bool = False
|
||||
|
||||
|
||||
def invalidate_vpn_cache() -> None:
|
||||
global _cache, _cache_at, _cache_failed
|
||||
_cache = None
|
||||
_cache_at = 0.0
|
||||
_cache_failed = False
|
||||
|
||||
|
||||
def fetch_vpn_settings(*, force: bool = False) -> VpnSettings | None:
|
||||
"""Return VPN settings from CA, or None if unreachable / unset."""
|
||||
global _cache, _cache_at, _cache_failed
|
||||
|
||||
now = time.monotonic()
|
||||
if (
|
||||
not force
|
||||
and _cache is not None
|
||||
and (now - _cache_at) < CACHE_TTL
|
||||
and not _cache_failed
|
||||
):
|
||||
return _cache
|
||||
|
||||
try:
|
||||
with httpx.Client(timeout=10.0) as client:
|
||||
response = client.get(
|
||||
f"{CA_API_URL}/internal/vpn",
|
||||
headers={"X-Internal-Token": INTERNAL_TOKEN},
|
||||
)
|
||||
response.raise_for_status()
|
||||
raw: dict[str, Any] = response.json()
|
||||
settings = VpnSettings.model_validate(raw)
|
||||
_cache = settings
|
||||
_cache_at = now
|
||||
_cache_failed = False
|
||||
return settings
|
||||
except Exception:
|
||||
logger.exception("Failed to load VPN settings from CA")
|
||||
_cache_failed = True
|
||||
_cache_at = now
|
||||
# Keep stale cache if present
|
||||
return _cache
|
||||
|
||||
|
||||
def proxy_fingerprint(settings: VpnSettings | None = None) -> str:
|
||||
cfg = settings if settings is not None else fetch_vpn_settings()
|
||||
if cfg is None or not cfg.proxies_source("telegram"):
|
||||
# Env fallback fingerprint
|
||||
from workers.sources.telegram_settings import get_env_proxy
|
||||
|
||||
env = get_env_proxy()
|
||||
return f"env:{env!r}"
|
||||
if cfg.mode == "subscription":
|
||||
return f"sub:{CP_VPN_HOST}:{CP_VPN_PORT}:{cfg.subscription_url}"
|
||||
return (
|
||||
f"{cfg.mode}:{cfg.host}:{cfg.port}:"
|
||||
f"{cfg.username or ''}:{'*' if cfg.password else ''}"
|
||||
)
|
||||
|
||||
|
||||
def telethon_proxy_tuple(settings: VpnSettings | None = None) -> tuple | None:
|
||||
"""Telethon-compatible proxy tuple for telegram traffic, or None."""
|
||||
cfg = settings if settings is not None else fetch_vpn_settings()
|
||||
if cfg is not None and cfg.proxies_source("telegram"):
|
||||
if cfg.mode == "subscription":
|
||||
return ("socks5", CP_VPN_HOST, CP_VPN_PORT)
|
||||
if cfg.mode in ("socks5", "http") and cfg.host and cfg.port:
|
||||
if cfg.username:
|
||||
return (
|
||||
cfg.mode,
|
||||
cfg.host,
|
||||
int(cfg.port),
|
||||
True,
|
||||
cfg.username,
|
||||
cfg.password or "",
|
||||
)
|
||||
return (cfg.mode, cfg.host, int(cfg.port))
|
||||
return None
|
||||
|
||||
from workers.sources.telegram_settings import get_env_proxy
|
||||
|
||||
return get_env_proxy()
|
||||
|
||||
|
||||
def http_proxy_url(source_type: str, settings: VpnSettings | None = None) -> str | None:
|
||||
"""HTTP(S)/SOCKS proxy URL for httpx / browsers, or None."""
|
||||
cfg = settings if settings is not None else fetch_vpn_settings()
|
||||
if cfg is None or not cfg.proxies_source(source_type):
|
||||
return None
|
||||
if cfg.mode == "subscription":
|
||||
return f"socks5://{CP_VPN_HOST}:{CP_VPN_PORT}"
|
||||
if cfg.mode in ("socks5", "http") and cfg.host and cfg.port:
|
||||
auth = ""
|
||||
if cfg.username:
|
||||
from urllib.parse import quote
|
||||
|
||||
user = quote(cfg.username, safe="")
|
||||
password = quote(cfg.password or "", safe="")
|
||||
auth = f"{user}:{password}@"
|
||||
return f"{cfg.mode}://{auth}{cfg.host}:{int(cfg.port)}"
|
||||
return None
|
||||
Reference in New Issue
Block a user