From 5811ecb134db84959a319871b45f2f309bbc0d79 Mon Sep 17 00:00:00 2001 From: gitrusprus Date: Sun, 16 Aug 2026 23:41:33 +0300 Subject: [PATCH] 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 --- .env.example | 5 +- AGENTS.md | 2 + centers/analytics/api/app/main.py | 6 +- centers/analytics/api/app/models.py | 19 ++ centers/analytics/api/app/routers/internal.py | 17 +- centers/analytics/api/app/routers/vpn.py | 32 ++ centers/analytics/api/app/schemas.py | 40 +++ centers/analytics/api/app/seed.py | 5 + centers/analytics/api/app/services/vpn.py | 137 +++++++++ centers/analytics/frontend/src/api/admin.ts | 14 + .../frontend/src/layouts/AppLayout.vue | 1 + .../analytics/frontend/src/router/index.ts | 7 + centers/analytics/frontend/src/types/admin.ts | 30 ++ .../analytics/frontend/src/views/VpnView.vue | 278 ++++++++++++++++++ centers/parsing/vpn/Dockerfile | 31 ++ centers/parsing/vpn/controller.py | 204 +++++++++++++ centers/parsing/workers/worker.py | 23 +- .../workers/adapters/crawl4ai_adapter.py | 20 +- .../parsing/workers/workers/adapters/viina.py | 9 +- .../workers/sources/telegram_session.py | 30 +- .../workers/sources/telegram_settings.py | 12 +- .../workers/workers/sources/vpn_config.py | 124 ++++++++ contracts/vpn.py | 71 +++++ docker-compose.yml | 29 +- docs/local-dev.md | 18 +- 25 files changed, 1146 insertions(+), 18 deletions(-) create mode 100644 centers/analytics/api/app/routers/vpn.py create mode 100644 centers/analytics/api/app/services/vpn.py create mode 100644 centers/analytics/frontend/src/views/VpnView.vue create mode 100644 centers/parsing/vpn/Dockerfile create mode 100644 centers/parsing/vpn/controller.py create mode 100644 centers/parsing/workers/workers/sources/vpn_config.py create mode 100644 contracts/vpn.py diff --git a/.env.example b/.env.example index 5e346c6..bc32015 100644 --- a/.env.example +++ b/.env.example @@ -3,10 +3,13 @@ TELEGRAM_API_ID=12345678 TELEGRAM_API_HASH=your_api_hash_here TELEGRAM_SESSION_PATH=/data/telegram.session -# Optional Telegram proxy +# Optional Telegram proxy fallback (when VPN tab is off) +# Prefer Admin UI → VPN (subscription / socks5 / http). # TELEGRAM_PROXY_TYPE=socks5 # TELEGRAM_PROXY_HOST=127.0.0.1 # TELEGRAM_PROXY_PORT=1080 +# TELEGRAM_PROXY_USER= +# TELEGRAM_PROXY_PASS= # Platform internals INTERNAL_TOKEN=dev-internal-token diff --git a/AGENTS.md b/AGENTS.md index fcd22ee..3029934 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -44,6 +44,8 @@ TOKEN=$(curl -s -X POST http://localhost:8080/admin/auth/login \ -d '{"username":"admin","password":"change-me"}' | jq -r .access_token) ``` +VPN (subscription / SOCKS): Admin UI → **VPN**, or `PUT /admin/vpn`. Workers and `cp-vpn` read `GET /internal/vpn`. + Create Telegram parser: ```bash diff --git a/centers/analytics/api/app/main.py b/centers/analytics/api/app/main.py index b466d49..1e79848 100644 --- a/centers/analytics/api/app/main.py +++ b/centers/analytics/api/app/main.py @@ -13,8 +13,8 @@ for _candidate in (_HERE.parent, *_HERE.parents): break from .database import Base, engine, get_db -from .routers import admin, auth, internal, map, objects, parse_channels, parser_profiles, v1 -from .seed import seed_objects, seed_test_consumer +from .routers import admin, auth, internal, map, objects, parse_channels, parser_profiles, v1, vpn +from .seed import seed_objects, seed_test_consumer, seed_vpn_settings from .services.migrations import migrate_schema from .services.scheduler import start_scheduler from .storage import ensure_upload_dir @@ -30,6 +30,7 @@ async def lifespan(_: FastAPI): try: seed_objects(db) seed_test_consumer(db) + seed_vpn_settings(db) finally: db.close() yield @@ -53,4 +54,5 @@ app.include_router(auth.router) app.include_router(admin.router) app.include_router(parser_profiles.router) app.include_router(parse_channels.router) +app.include_router(vpn.router) app.include_router(v1.router) diff --git a/centers/analytics/api/app/models.py b/centers/analytics/api/app/models.py index 8ed45ae..d8c6bd6 100644 --- a/centers/analytics/api/app/models.py +++ b/centers/analytics/api/app/models.py @@ -190,3 +190,22 @@ class ConsumerFilter(Base): date_from: Mapped[datetime | None] = mapped_column(DateTime(timezone=True), nullable=True) consumer: Mapped["Consumer"] = relationship(back_populates="filter") + + +class VpnSettingsRow(Base): + """Singleton VPN / proxy settings (id=1).""" + + __tablename__ = "vpn_settings" + + id: Mapped[int] = mapped_column(Integer, primary_key=True) + enabled: Mapped[bool] = mapped_column(Boolean, default=False, nullable=False) + mode: Mapped[str] = mapped_column(String(32), default="subscription", nullable=False) + subscription_url: Mapped[str | None] = mapped_column(Text, nullable=True) + subscription_interval_seconds: Mapped[int] = mapped_column( + Integer, default=3600, nullable=False + ) + host: Mapped[str | None] = mapped_column(String(255), nullable=True) + port: Mapped[int | None] = mapped_column(Integer, nullable=True) + username: Mapped[str | None] = mapped_column(String(255), nullable=True) + password: Mapped[str | None] = mapped_column(String(255), nullable=True) + proxied_source_types: Mapped[list | None] = mapped_column(JSON, nullable=True) diff --git a/centers/analytics/api/app/routers/internal.py b/centers/analytics/api/app/routers/internal.py index a1a00dd..7afe57d 100644 --- a/centers/analytics/api/app/routers/internal.py +++ b/centers/analytics/api/app/routers/internal.py @@ -6,9 +6,15 @@ from sqlalchemy.orm import Session from ..database import get_db from ..deps import verify_internal_token from ..models import ParseJob -from ..schemas import IngestRequest, IngestResponse, ListenerSubscription +from ..schemas import ( + IngestRequest, + IngestResponse, + ListenerSubscription, + VpnSettingsInternalRead, +) from ..services.ingest import ingest_events from ..services.jobs import resolve_job_source_config +from ..services.vpn import ensure_vpn_settings, to_internal_read router = APIRouter(prefix="/internal", tags=["internal"]) @@ -85,3 +91,12 @@ def listener_subscriptions( ) db.commit() return result + + +@router.get("/vpn", response_model=VpnSettingsInternalRead) +def internal_vpn_settings( + _: None = Depends(verify_internal_token), + db: Session = Depends(get_db), +): + row = ensure_vpn_settings(db) + return to_internal_read(row) diff --git a/centers/analytics/api/app/routers/vpn.py b/centers/analytics/api/app/routers/vpn.py new file mode 100644 index 0000000..3e7838c --- /dev/null +++ b/centers/analytics/api/app/routers/vpn.py @@ -0,0 +1,32 @@ +"""Admin VPN / proxy settings.""" + +from __future__ import annotations + +from fastapi import APIRouter, Depends, HTTPException +from sqlalchemy.orm import Session + +from ..database import get_db +from ..deps import verify_admin +from ..schemas import VpnSettingsAdminRead, VpnSettingsUpdate +from ..services.vpn import apply_vpn_update, ensure_vpn_settings, to_admin_read + +router = APIRouter( + prefix="/admin/vpn", + tags=["vpn"], + dependencies=[Depends(verify_admin)], +) + + +@router.get("", response_model=VpnSettingsAdminRead) +def get_vpn_settings(db: Session = Depends(get_db)): + row = ensure_vpn_settings(db) + return to_admin_read(row) + + +@router.put("", response_model=VpnSettingsAdminRead) +def put_vpn_settings(payload: VpnSettingsUpdate, db: Session = Depends(get_db)): + try: + row = apply_vpn_update(db, payload) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + return to_admin_read(row) diff --git a/centers/analytics/api/app/schemas.py b/centers/analytics/api/app/schemas.py index dc20b19..1896388 100644 --- a/centers/analytics/api/app/schemas.py +++ b/centers/analytics/api/app/schemas.py @@ -300,3 +300,43 @@ class TimelinePoint(BaseModel): class TopItem(BaseModel): name: str count: int + + +class VpnSettingsUpdate(BaseModel): + enabled: bool | None = None + mode: Literal["subscription", "socks5", "http"] | None = None + subscription_url: str | None = None + clear_subscription_url: bool | None = None + subscription_interval_seconds: int | None = Field(default=None, ge=60, le=86400) + host: str | None = None + port: int | None = Field(default=None, ge=1, le=65535) + username: str | None = None + password: str | None = None + clear_password: bool | None = None + proxied_source_types: list[str] | None = None + + +class VpnSettingsAdminRead(BaseModel): + enabled: bool + mode: Literal["subscription", "socks5", "http"] + subscription_url_set: bool + subscription_url_hint: str | None = None + subscription_interval_seconds: int + host: str | None = None + port: int | None = None + username: str | None = None + password_set: bool + proxied_source_types: list[str] + available_source_types: list[str] + + +class VpnSettingsInternalRead(BaseModel): + enabled: bool + mode: Literal["subscription", "socks5", "http"] + subscription_url: str | None = None + subscription_interval_seconds: int + host: str | None = None + port: int | None = None + username: str | None = None + password: str | None = None + proxied_source_types: list[str] diff --git a/centers/analytics/api/app/seed.py b/centers/analytics/api/app/seed.py index ce17ff1..7f7d45c 100644 --- a/centers/analytics/api/app/seed.py +++ b/centers/analytics/api/app/seed.py @@ -4,6 +4,7 @@ from sqlalchemy.orm import Session from .models import Consumer, ConsumerFilter from .services.filtering import hash_api_key +from .services.vpn import ensure_vpn_settings def seed_objects(db: Session) -> None: @@ -24,3 +25,7 @@ def seed_test_consumer(db: Session) -> None: db.flush() db.add(ConsumerFilter(consumer_id=consumer.id, regions=None, topics=None)) db.commit() + + +def seed_vpn_settings(db: Session) -> None: + ensure_vpn_settings(db) diff --git a/centers/analytics/api/app/services/vpn.py b/centers/analytics/api/app/services/vpn.py new file mode 100644 index 0000000..b0b3ebd --- /dev/null +++ b/centers/analytics/api/app/services/vpn.py @@ -0,0 +1,137 @@ +"""VPN settings persistence helpers.""" + +from __future__ import annotations + +from contracts.queues import known_source_types +from contracts.vpn import VpnSettings +from sqlalchemy.orm import Session + +from ..models import VpnSettingsRow +from ..schemas import VpnSettingsAdminRead, VpnSettingsInternalRead, VpnSettingsUpdate + +SINGLETON_ID = 1 + + +def _hint_secret(value: str | None) -> str | None: + if not value: + return None + if len(value) <= 4: + return "****" + return f"…{value[-4:]}" + + +def ensure_vpn_settings(db: Session) -> VpnSettingsRow: + row = db.query(VpnSettingsRow).filter(VpnSettingsRow.id == SINGLETON_ID).first() + if row is not None: + return row + row = VpnSettingsRow( + id=SINGLETON_ID, + enabled=False, + mode="subscription", + subscription_interval_seconds=3600, + proxied_source_types=[], + ) + db.add(row) + db.commit() + db.refresh(row) + return row + + +def row_to_contract(row: VpnSettingsRow) -> VpnSettings: + return VpnSettings.model_validate( + { + "enabled": row.enabled, + "mode": row.mode, + "subscription_url": row.subscription_url, + "subscription_interval_seconds": row.subscription_interval_seconds, + "host": row.host, + "port": row.port, + "username": row.username, + "password": row.password, + "proxied_source_types": row.proxied_source_types or [], + } + ) + + +def to_admin_read(row: VpnSettingsRow) -> VpnSettingsAdminRead: + return VpnSettingsAdminRead( + enabled=row.enabled, + mode=row.mode, # type: ignore[arg-type] + subscription_url_set=bool(row.subscription_url), + subscription_url_hint=_hint_secret(row.subscription_url), + subscription_interval_seconds=row.subscription_interval_seconds, + host=row.host, + port=row.port, + username=row.username, + password_set=bool(row.password), + proxied_source_types=list(row.proxied_source_types or []), + available_source_types=known_source_types(), + ) + + +def to_internal_read(row: VpnSettingsRow) -> VpnSettingsInternalRead: + return VpnSettingsInternalRead( + enabled=row.enabled, + mode=row.mode, # type: ignore[arg-type] + subscription_url=row.subscription_url, + subscription_interval_seconds=row.subscription_interval_seconds, + host=row.host, + port=row.port, + username=row.username, + password=row.password, + proxied_source_types=list(row.proxied_source_types or []), + ) + + +def apply_vpn_update(db: Session, payload: VpnSettingsUpdate) -> VpnSettingsRow: + row = ensure_vpn_settings(db) + data = { + "enabled": row.enabled if payload.enabled is None else payload.enabled, + "mode": row.mode if payload.mode is None else payload.mode, + "subscription_url": row.subscription_url, + "subscription_interval_seconds": ( + row.subscription_interval_seconds + if payload.subscription_interval_seconds is None + else payload.subscription_interval_seconds + ), + "host": row.host if payload.host is None else (payload.host.strip() or None), + "port": row.port if payload.port is None else payload.port, + "username": ( + row.username + if payload.username is None + else (payload.username.strip() or None) + ), + "password": row.password, + "proxied_source_types": ( + list(row.proxied_source_types or []) + if payload.proxied_source_types is None + else payload.proxied_source_types + ), + } + + if payload.clear_subscription_url: + data["subscription_url"] = None + elif payload.subscription_url is not None: + stripped = payload.subscription_url.strip() + data["subscription_url"] = stripped or None + + if payload.clear_password: + data["password"] = None + elif payload.password is not None: + stripped = payload.password.strip() + data["password"] = stripped or None + + validated = VpnSettings.model_validate(data) + + row.enabled = validated.enabled + row.mode = validated.mode + row.subscription_url = validated.subscription_url + row.subscription_interval_seconds = validated.subscription_interval_seconds + row.host = validated.host + row.port = validated.port + row.username = validated.username + row.password = validated.password + row.proxied_source_types = validated.proxied_source_types + db.commit() + db.refresh(row) + return row diff --git a/centers/analytics/frontend/src/api/admin.ts b/centers/analytics/frontend/src/api/admin.ts index d7a9b3a..8ce28fd 100644 --- a/centers/analytics/frontend/src/api/admin.ts +++ b/centers/analytics/frontend/src/api/admin.ts @@ -14,6 +14,8 @@ import type { ParseJobUpdate, TimelinePoint, TopItem, + VpnSettings, + VpnSettingsUpdate, } from "../types/admin"; function buildQuery(params: Record): string { @@ -157,3 +159,15 @@ export async function testDistribution(apiKey: string, limit = 5): Promise; } + +export function fetchVpnSettings(): Promise { + return request("/vpn", undefined, ADMIN_BASE); +} + +export function updateVpnSettings(payload: VpnSettingsUpdate): Promise { + return request( + "/vpn", + { method: "PUT", body: JSON.stringify(payload) }, + ADMIN_BASE, + ); +} diff --git a/centers/analytics/frontend/src/layouts/AppLayout.vue b/centers/analytics/frontend/src/layouts/AppLayout.vue index 6763a16..b8d9a20 100644 --- a/centers/analytics/frontend/src/layouts/AppLayout.vue +++ b/centers/analytics/frontend/src/layouts/AppLayout.vue @@ -13,6 +13,7 @@ const adminNavItems = [ { to: "/channels", label: "Каналы" }, { to: "/parser-profiles", label: "Профили" }, { to: "/parsers", label: "Парсеры" }, + { to: "/vpn", label: "VPN" }, { to: "/events", label: "События" }, { to: "/analytics", label: "Аналитика" }, { to: "/consumers", label: "ПИ" }, diff --git a/centers/analytics/frontend/src/router/index.ts b/centers/analytics/frontend/src/router/index.ts index 4327460..0ac481a 100644 --- a/centers/analytics/frontend/src/router/index.ts +++ b/centers/analytics/frontend/src/router/index.ts @@ -10,6 +10,7 @@ import LoginView from "../views/LoginView.vue"; import MapViewPage from "../views/MapViewPage.vue"; import ParserProfilesView from "../views/ParserProfilesView.vue"; import ParsersView from "../views/ParsersView.vue"; +import VpnView from "../views/VpnView.vue"; const router = createRouter({ history: createWebHistory(), @@ -43,6 +44,12 @@ const router = createRouter({ component: ParsersView, meta: { requiresAuth: true }, }, + { + path: "vpn", + name: "vpn", + component: VpnView, + meta: { requiresAuth: true }, + }, { path: "events", name: "events", diff --git a/centers/analytics/frontend/src/types/admin.ts b/centers/analytics/frontend/src/types/admin.ts index 7334f83..05cc647 100644 --- a/centers/analytics/frontend/src/types/admin.ts +++ b/centers/analytics/frontend/src/types/admin.ts @@ -206,3 +206,33 @@ export interface TopItem { name: string; count: number; } + +export type VpnMode = "subscription" | "socks5" | "http"; + +export interface VpnSettings { + enabled: boolean; + mode: VpnMode; + subscription_url_set: boolean; + subscription_url_hint: string | null; + subscription_interval_seconds: number; + host: string | null; + port: number | null; + username: string | null; + password_set: boolean; + proxied_source_types: string[]; + available_source_types: string[]; +} + +export interface VpnSettingsUpdate { + enabled?: boolean; + mode?: VpnMode; + subscription_url?: string | null; + clear_subscription_url?: boolean; + subscription_interval_seconds?: number; + host?: string | null; + port?: number | null; + username?: string | null; + password?: string | null; + clear_password?: boolean; + proxied_source_types?: string[]; +} diff --git a/centers/analytics/frontend/src/views/VpnView.vue b/centers/analytics/frontend/src/views/VpnView.vue new file mode 100644 index 0000000..2569180 --- /dev/null +++ b/centers/analytics/frontend/src/views/VpnView.vue @@ -0,0 +1,278 @@ + + + + + diff --git a/centers/parsing/vpn/Dockerfile b/centers/parsing/vpn/Dockerfile new file mode 100644 index 0000000..1dee542 --- /dev/null +++ b/centers/parsing/vpn/Dockerfile @@ -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"] diff --git a/centers/parsing/vpn/controller.py b/centers/parsing/vpn/controller.py new file mode 100644 index 0000000..7fbcf98 --- /dev/null +++ b/centers/parsing/vpn/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() diff --git a/centers/parsing/workers/worker.py b/centers/parsing/workers/worker.py index 426dae0..8f17d0d 100644 --- a/centers/parsing/workers/worker.py +++ b/centers/parsing/workers/worker.py @@ -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)) diff --git a/centers/parsing/workers/workers/adapters/crawl4ai_adapter.py b/centers/parsing/workers/workers/adapters/crawl4ai_adapter.py index e1a9ac8..315b1c7 100644 --- a/centers/parsing/workers/workers/adapters/crawl4ai_adapter.py +++ b/centers/parsing/workers/workers/adapters/crawl4ai_adapter.py @@ -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": diff --git a/centers/parsing/workers/workers/adapters/viina.py b/centers/parsing/workers/workers/adapters/viina.py index 88dd32d..178ac3e 100644 --- a/centers/parsing/workers/workers/adapters/viina.py +++ b/centers/parsing/workers/workers/adapters/viina.py @@ -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"}, diff --git a/centers/parsing/workers/workers/sources/telegram_session.py b/centers/parsing/workers/workers/sources/telegram_session.py index 1632144..cb3828e 100644 --- a/centers/parsing/workers/workers/sources/telegram_session.py +++ b/centers/parsing/workers/workers/sources/telegram_session.py @@ -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 diff --git a/centers/parsing/workers/workers/sources/telegram_settings.py b/centers/parsing/workers/workers/sources/telegram_settings.py index cde5e0a..40307d3 100644 --- a/centers/parsing/workers/workers/sources/telegram_settings.py +++ b/centers/parsing/workers/workers/sources/telegram_settings.py @@ -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: diff --git a/centers/parsing/workers/workers/sources/vpn_config.py b/centers/parsing/workers/workers/sources/vpn_config.py new file mode 100644 index 0000000..195caba --- /dev/null +++ b/centers/parsing/workers/workers/sources/vpn_config.py @@ -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 diff --git a/contracts/vpn.py b/contracts/vpn.py new file mode 100644 index 0000000..2248112 --- /dev/null +++ b/contracts/vpn.py @@ -0,0 +1,71 @@ +"""VPN / proxy settings shared by CA admin UI and CP workers.""" + +from __future__ import annotations + +from typing import Literal + +from pydantic import BaseModel, Field, field_validator, model_validator + +from contracts.queues import known_source_types + +VpnMode = Literal["subscription", "socks5", "http"] + +VPN_MODES: tuple[str, ...] = ("subscription", "socks5", "http") + + +class VpnSettings(BaseModel): + enabled: bool = False + mode: VpnMode = "subscription" + subscription_url: str | None = None + subscription_interval_seconds: int = Field(default=3600, ge=60, le=86400) + host: str | None = None + port: int | None = Field(default=None, ge=1, le=65535) + username: str | None = None + password: str | None = None + proxied_source_types: list[str] = Field(default_factory=list) + + @field_validator("subscription_url", "host", "username", "password", mode="before") + @classmethod + def empty_to_none(cls, value: object) -> object: + if value is None: + return None + if isinstance(value, str): + stripped = value.strip() + return stripped or None + return value + + @field_validator("proxied_source_types") + @classmethod + def validate_source_types(cls, value: list[str]) -> list[str]: + known = set(known_source_types()) + cleaned: list[str] = [] + seen: set[str] = set() + for raw in value or []: + item = str(raw).strip() + if not item or item in seen: + continue + if item not in known: + raise ValueError( + f"Unknown source_type {item!r}. " + f"Allowed: {', '.join(sorted(known))}" + ) + seen.add(item) + cleaned.append(item) + return cleaned + + @model_validator(mode="after") + def require_fields_when_enabled(self) -> "VpnSettings": + if not self.enabled: + return self + if self.mode == "subscription": + if not self.subscription_url: + raise ValueError("subscription_url required when mode=subscription") + elif self.mode in ("socks5", "http"): + if not self.host: + raise ValueError(f"host required when mode={self.mode}") + if self.port is None: + raise ValueError(f"port required when mode={self.mode}") + return self + + def proxies_source(self, source_type: str) -> bool: + return self.enabled and source_type in self.proxied_source_types diff --git a/docker-compose.yml b/docker-compose.yml index 4089416..2253a5d 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -49,6 +49,17 @@ services: - ca-api restart: unless-stopped + cp-vpn: + build: ./centers/parsing/vpn + environment: + CA_API_URL: http://ca-api:8000 + INTERNAL_TOKEN: ${INTERNAL_TOKEN:-dev-internal-token} + VPN_POLL_SECONDS: "15" + VPN_SOCKS_PORT: "1080" + depends_on: + - ca-api + restart: unless-stopped + cp-workers: build: context: . @@ -62,14 +73,17 @@ services: TELEGRAM_SESSION_PATH: /data/telegram.session TELEGRAM_LISTENER_ENABLED: "true" TELEGRAM_LISTENER_REFRESH_SECONDS: "60" - # Host xray is 127.0.0.1:10808; socat on host forwards 0.0.0.0:11080 → that SOCKS + # Env fallback when VPN tab is off (host Happ/xray; socat 11080→10808). + # Overrides .env 127.0.0.1 which is wrong inside the container. TELEGRAM_PROXY_HOST: host.docker.internal TELEGRAM_PROXY_PORT: "11080" + CP_VPN_HOST: cp-vpn + CP_VPN_PORT: "1080" ENABLED_ADAPTERS: telegram WORKER_FAMILIES: telegram CA_API_URL: http://ca-api:8000 REDIS_URL: redis://redis:6379/0 - INTERNAL_TOKEN: dev-internal-token + INTERNAL_TOKEN: ${INTERNAL_TOKEN:-dev-internal-token} extra_hosts: - "host.docker.internal:host-gateway" volumes: @@ -77,6 +91,7 @@ services: depends_on: - redis - ca-api + - cp-vpn restart: unless-stopped cp-workers-web: @@ -92,12 +107,15 @@ services: TELEGRAM_LISTENER_ENABLED: "false" ENABLED_ADAPTERS: crawl4ai WORKER_FAMILIES: web + CP_VPN_HOST: cp-vpn + CP_VPN_PORT: "1080" CA_API_URL: http://ca-api:8000 REDIS_URL: redis://redis:6379/0 - INTERNAL_TOKEN: dev-internal-token + INTERNAL_TOKEN: ${INTERNAL_TOKEN:-dev-internal-token} depends_on: - redis - ca-api + - cp-vpn restart: unless-stopped cp-workers-nlp: @@ -111,12 +129,15 @@ services: TELEGRAM_LISTENER_ENABLED: "false" ENABLED_ADAPTERS: viina WORKER_FAMILIES: nlp + CP_VPN_HOST: cp-vpn + CP_VPN_PORT: "1080" CA_API_URL: http://ca-api:8000 REDIS_URL: redis://redis:6379/0 - INTERNAL_TOKEN: dev-internal-token + INTERNAL_TOKEN: ${INTERNAL_TOKEN:-dev-internal-token} depends_on: - redis - ca-api + - cp-vpn restart: unless-stopped volumes: diff --git a/docs/local-dev.md b/docs/local-dev.md index 0f1ba78..e18a633 100644 --- a/docs/local-dev.md +++ b/docs/local-dev.md @@ -32,6 +32,7 @@ Health: `curl http://localhost:8080/api/health` | `cp-workers` | — | Telegram | | `cp-workers-web` | — | Crawl4AI | | `cp-workers-nlp` | — | VIINA | +| `cp-vpn` | внутренний 1080 | SOCKS из subscription (вкладка VPN) | Пересборка одного воркера: @@ -52,7 +53,20 @@ docker compose config - Файл **не** коммитить (`.gitignore`) - API id/hash в `.env` должны совпадать с теми, под которыми создавалась сессия -### Создать / обновить сессию (локальный прокси) +## VPN / прокси (вкладка VPN) + +Основной путь: UI → **VPN** (`/vpn`). + +1. Включите «Прокси включён». +2. Способ **Subscription URL** — вставьте ссылку подписки (не коммитьте её в git). +3. Отметьте источники (`telegram`, при необходимости `crawl4ai` / `viina`). +4. Сохраните. Сервис `cp-vpn` (mihomo) подтянет подписку и отдаст SOCKS на `cp-vpn:1080`; воркеры читают `/internal/vpn`. + +Режимы **SOCKS5** / **HTTP** задают host:port напрямую (без `cp-vpn` для Telegram). + +### Fallback: локальный Happ/xray (без вкладки VPN) + +Если VPN в UI выключен, воркеры используют `TELEGRAM_PROXY_*` из `.env`: 1. Остановите Telegram-воркер, чтобы не делить session-файл: `docker compose stop cp-workers` @@ -61,7 +75,7 @@ docker compose config ```bash socat TCP-LISTEN:11080,bind=0.0.0.0,fork,reuseaddr TCP:127.0.0.1:10808 & ``` - В compose у `cp-workers`: `TELEGRAM_PROXY_HOST=host.docker.internal`, порт `11080`. + Для auth из контейнера с host-сетью: `TELEGRAM_PROXY_HOST=127.0.0.1`. 4. Авторизация (интерактивно, код из Telegram): ```bash docker run --rm -it --network host \