Encrypted proxy killswitch + TUN optimization #1

Open
zenaku wants to merge 54 commits from feature/merge-encrypted-proxy into master
41 changed files with 2165 additions and 276 deletions

2
.gitignore vendored
View file

@ -4,3 +4,5 @@ prototype_client.py
.venv
__pycache__
dist
.env
docs.md

163
README.md
View file

@ -1,34 +1,163 @@
# sp-hydra-veil-core
# Hydra-Veil — June Major Update
The `sp-hydra-veil-core` library exposes core logic to higher-level components.
## Expanded Killswitch
## Build Instructions
`hydraveil-killswitch.sh` was added to the installer: a wrapper around `nftables` that blocks all traffic not routed through the tunnel. It defaults to `tun0`, but supports a configurable interface IP that the operator will send when creating the profile. The IP is currently hardcoded; it will be removed in the next commit.
### Presumptions
## One Connection at a Time
* Your system is configured to use the Simplified Privacy package registry [1].
First in, first out. When a new connection is established, the previous one receives a signal (`SIGUSR1`) and closes cleanly before the new one takes over.
### Prerequisites
## Unexpected Shutdowns
* `build`
* `twine`
If the terminal is closed or the machine is powered off, the process receives `SIGHUP` and performs a full disconnect: stops `sing-box`, clears state, and releases the network. No orphan processes.
To install them, activate your `venv`, if necessary, and run:
## Fewer Network Requests
```bash
pip install build twine
Unnecessary calls made when establishing and closing connections have been removed.
---
## Repositories
| Repo | Role |
|---|---|
| `core` | Business logic, controllers, models |
| `cli` | Command-line interface |
| `essentials` | Network modules (Tor, WireGuard, proxies) |
| `installer` | Installation scripts and killswitch |
---
---
# Hydra-Veil — Gran Actualización de Junio
## Killswitch ampliado
Se añadió `hydraveil-killswitch.sh` al installer: un wrapper sobre `nftables` que bloquea todo el tráfico que no pase por el túnel. Por defecto usa `tun0`, pero soporta una IP de interfaz configurable que el operador enviará al crear el perfil. Por ahora la IP está hardcodeada; se eliminará en el próximo commit.
## Una conexión a la vez
Primero en entrar, primero en salir. Cuando se establece una nueva conexión, la anterior recibe una señal (`SIGUSR1`) y se cierra limpiamente antes de que la nueva tome el control.
## Cierres inesperados
Si se cierra la terminal o se apaga el equipo, el proceso recibe `SIGHUP` y ejecuta la desconexión completa: para `sing-box`, limpia el estado y libera la red. Sin procesos huérfanos.
## Menos solicitudes de red
Se eliminaron llamadas innecesarias que se hacían al establecer y cerrar conexiones.
---
## Repositorios
| Repo | Rol |
|---|---|
| `core` | Lógica de negocio, controladores, modelos |
| `cli` | Interfaz de línea de comandos |
| `essentials` | Módulos de red (Tor, WireGuard, proxies) |
| `installer` | Scripts de instalación y killswitch |
---
---
## Uso rápido — Desde Core (innecesario, usar CLI)
## 1. Create billing code / Crear billing code
```python
python3 -c "
from core.services.WebServiceApiService import WebServiceApiService
subscription = WebServiceApiService.post_subscription(2, operator_id=<OPERATOR_ID>)
print('billing code:', subscription.billing_code)
"
```
### Build Steps
> Copy the printed `billing_code` — you'll need it in the connect step.
> Copiá el `billing_code` impreso, lo vas a necesitar en el paso de conexión.
* Activate your `venv`, if necessary, and run:
---
consultar estado de invoive
python3 -c "
from core.services.WebServiceApiService import WebServiceApiService
invoice = WebServiceApiService.get_invoice('BILLING_CODE_HERE')
print('status:', invoice.status)
"
## 2. Patch subscription / Parchear suscripción
```bash
rm -rf ./dist/*
python3 -m build
python3 -m twine upload --repository forgejo ./dist/*
sudo docker compose exec laravel php artisan tinker --execute="
\$sub = \App\Models\Subscription::orderBy('id', 'desc')->first();
\$sub->expires_at = '2027-12-31 23:59:59';
\$sub->duration = 720;
\$sub->payment_reference = 'bypass_' . uniqid();
\$sub->save();
echo 'OK' . PHP_EOL;
"
```
---
[1] https://forgejo.org/docs/v14.0/user/packages/pypi/#configuring-the-package-registry
## 3. Connect / Conectar
### VLESS
```python
python3 -c "
from core.services.WebServiceApiService import WebServiceApiService
from core.controllers.encrypted_proxy.VlessController import VlessController
from core.observers.EncryptedProxyObserver import EncryptedProxyObserver
observer = EncryptedProxyObserver()
observer.subscribe('connected', lambda e: print('connected:', e.subject))
observer.subscribe('error', lambda e: print('error:', e.subject))
observer.subscribe('disconnected', lambda e: print('disconnected:', e.subject))
session = WebServiceApiService.post_operator_proxy('MZKD-QWWI-1TS9-KHMD', 3, 'vless')
print('session:', session)
controller = VlessController(1080)
controller.enable(session.links[0], session.username, observer)
"
```
### Hysteria2
```python
python3 -c "
from core.services.WebServiceApiService import WebServiceApiService
from core.controllers.encrypted_proxy.HysteriaController import HysteriaController
from core.observers.EncryptedProxyObserver import EncryptedProxyObserver
observer = EncryptedProxyObserver()
observer.subscribe('connected', lambda e: print('connected:', e.subject))
observer.subscribe('error', lambda e: print('error:', e.subject))
observer.subscribe('disconnected', lambda e: print('disconnected:', e.subject))
session = WebServiceApiService.post_operator_proxy('MZKD-QWWI-1TS9-KHMD', 3, 'hysteria2')
print('session:', session)
controller = HysteriaController(1080)
controller.enable(session.username, session.password, session.operator_hysteria2_host, observer)
"
```
---
## 4. Disconnect / Desconectar
```python
python3 -c "
from core.controllers.encrypted_proxy.HysteriaController import HysteriaController
from core.observers.EncryptedProxyObserver import EncryptedProxyObserver
observer = EncryptedProxyObserver()
observer.subscribe('disconnected', lambda e: print('disconnected'))
controller = HysteriaController(1080)
controller.disable(observer)
"
```
sudo ln -sf /run/systemd/resolve/resolv.conf /etc/resolv.conf

View file

@ -9,52 +9,72 @@ class Constants:
DB_VERSION_THIS_APP_WANTS = 1
# ticketing group:
TICKET_API_BASE_URL: Final[str] = os.environ.get(
TICKET_API_BASE_URL: Final[str] = os.environ.get(
"TICKET_API_BASE_URL", "https://ticket.hydraveil.net"
)
SP_API_BASE_URL: Final[str] = os.environ.get(
'SP_API_BASE_URL', 'https://fake.simplifiedprivacy.org/api/v1'
)
PING_URL: Final[str] = os.environ.get(
'PING_URL', 'https://fake.simplifiedprivacy.org/api/v1//health'
)
CONNECTION_RETRY_INTERVAL: Final[int] = int(os.environ.get('CONNECTION_RETRY_INTERVAL', '5'))
MAX_CONNECTION_ATTEMPTS: Final[int] = int(os.environ.get('MAX_CONNECTION_ATTEMPTS', '2'))
HV_CLIENT_PATH: Final[str] = os.environ.get('HV_CLIENT_PATH')
HV_CLIENT_VERSION_NUMBER: Final[str] = os.environ.get('HV_CLIENT_VERSION_NUMBER')
SP_API_BASE_URL: Final[str] = os.environ.get('SP_API_BASE_URL', 'https://api.hydraveil.net/api/v1')
PING_URL: Final[str] = os.environ.get('PING_URL', 'https://api.hydraveil.net/api/v1/health')
# ── Paths base ────────────────────────────────────────────────────────────
HOME: Final[str] = os.path.expanduser('~')
SYSTEM_CONFIG_PATH: Final[str] = '/etc'
CONNECTION_RETRY_INTERVAL: Final[int] = int(os.environ.get('CONNECTION_RETRY_INTERVAL', '5'))
MAX_CONNECTION_ATTEMPTS: Final[int] = int(os.environ.get('MAX_CONNECTION_ATTEMPTS', '2'))
# ── XDG estándar ─────────────────────────────────────────────────────────
CACHE_HOME: Final[str] = os.environ.get('XDG_CACHE_HOME', os.path.join(HOME, '.cache'))
CONFIG_HOME: Final[str] = os.environ.get('XDG_CONFIG_HOME', os.path.join(HOME, '.config'))
DATA_HOME: Final[str] = os.environ.get('XDG_DATA_HOME', os.path.join(HOME, '.local/share'))
STATE_HOME: Final[str] = os.environ.get('XDG_STATE_HOME', os.path.join(HOME, '.local/state'))
HV_CLIENT_PATH: Final[str] = os.environ.get('HV_CLIENT_PATH')
HV_CLIENT_VERSION_NUMBER: Final[str] = os.environ.get('HV_CLIENT_VERSION_NUMBER')
HOME: Final[str] = os.path.expanduser('~')
SYSTEM_CONFIG_PATH: Final[str] = '/etc'
CACHE_HOME: Final[str] = os.environ.get('XDG_CACHE_HOME', os.path.join(HOME, '.cache'))
CONFIG_HOME: Final[str] = os.environ.get('XDG_CONFIG_HOME', os.path.join(HOME, '.config'))
DATA_HOME: Final[str] = os.environ.get('XDG_DATA_HOME', os.path.join(HOME, '.local/share'))
STATE_HOME: Final[str] = os.environ.get('XDG_STATE_HOME', os.path.join(HOME, '.local/state'))
HV_SYSTEM_CONFIG_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/hydra-veil'
HV_CACHE_HOME: Final[str] = f'{CACHE_HOME}/hydra-veil'
HV_CONFIG_HOME: Final[str] = f'{CONFIG_HOME}/hydra-veil'
HV_DATA_HOME: Final[str] = f'{DATA_HOME}/hydra-veil'
HV_STATE_HOME: Final[str] = f'{STATE_HOME}/hydra-veil'
# ── hydra-veil dirs ───────────────────────────────────────────────────────
HV_SYSTEM_CONFIG_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/hydra-veil'
HV_CACHE_HOME: Final[str] = f'{CACHE_HOME}/hydra-veil'
HV_CONFIG_HOME: Final[str] = f'{CONFIG_HOME}/hydra-veil'
HV_DATA_HOME: Final[str] = f'{DATA_HOME}/hydra-veil'
HV_STATE_HOME: Final[str] = f'{STATE_HOME}/hydra-veil'
HV_SYSTEM_PROFILE_CONFIG_PATH: Final[str] = f'{HV_SYSTEM_CONFIG_PATH}/profiles'
HV_PROFILE_CONFIG_HOME: Final[str] = f'{HV_CONFIG_HOME}/profiles'
HV_PROFILE_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/profiles'
HV_TICKETING_CONFIG_HOME: Final[str] = f'{HV_CONFIG_HOME}/ticketing'
HV_TICKETING_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/ticket_data'
HV_APPLICATION_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/applications'
HV_INCIDENT_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/incidents'
HV_RUNTIME_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/runtime'
HV_STORAGE_DATABASE_PATH: Final[str] = f'{HV_DATA_HOME}/storage.db'
HV_CAPABILITY_POLICY_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/apparmor.d/hydra-veil'
HV_PRIVILEGE_POLICY_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/sudoers.d/hydra-veil'
HV_SESSION_STATE_HOME: Final[str] = f'{HV_STATE_HOME}/sessions'
HV_TOR_STATE_HOME: Final[str] = f'{HV_STATE_HOME}/tor'
HV_PROFILE_CONFIG_HOME: Final[str] = f'{HV_CONFIG_HOME}/profiles'
HV_PROFILE_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/profiles'
# ── sing-box ──────────────────────────────────────────────────────────────
SINGBOX_WRAPPER: Final[str] = os.environ.get(
'SINGBOX_WRAPPER', '/usr/local/bin/hydraveil-singbox'
)
SINGBOX_BIN: Final[str] = os.environ.get(
'SINGBOX_BIN', '/usr/bin/sing-box'
)
SINGBOX_CONFIG_DIR: Final[str] = f'{HV_DATA_HOME}/configs'
SINGBOX_PID_FILE: Final[str] = f'{HV_RUNTIME_DATA_HOME}/singbox.pid'
SINGBOX_LOG_FILE: Final[str] = f'{HV_RUNTIME_DATA_HOME}/singbox.log'
SINGBOX_TUN_IF: Final[str] = os.environ.get('SINGBOX_TUN_IF', 'tun0')
SINGBOX_INTERNAL_SUBNET: Final[str] = os.environ.get('SINGBOX_INTERNAL_SUBNET', '172.19.0.0/30')
SINGBOX_INTERNAL_ADDR: Final[str] = os.environ.get('SINGBOX_INTERNAL_ADDR', '172.19.0.1/30')
# ticketing group:
HV_TICKETING_CONFIG_HOME: Final[str] = f"{HV_CONFIG_HOME}/ticketing"
HV_TICKETING_DATA_HOME: Final[str] = f"{HV_DATA_HOME}/ticket_data"
# ── killswitch / dns wrappers ─────────────────────────────────────────────
KILLSWITCH_WRAPPER: Final[str] = os.environ.get(
'KILLSWITCH_WRAPPER', '/usr/local/bin/hydraveil-killswitch'
)
RESOLVECTL_WRAPPER: Final[str] = os.environ.get(
'RESOLVECTL_WRAPPER', '/usr/local/bin/hydraveil-resolvectl'
)
HV_APPLICATION_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/applications'
HV_INCIDENT_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/incidents'
HV_RUNTIME_DATA_HOME: Final[str] = f'{HV_DATA_HOME}/runtime'
HV_STORAGE_DATABASE_PATH: Final[str] = f'{HV_DATA_HOME}/storage.db'
HV_CAPABILITY_POLICY_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/apparmor.d/hydra-veil'
HV_PRIVILEGE_POLICY_PATH: Final[str] = f'{SYSTEM_CONFIG_PATH}/sudoers.d/hydra-veil'
HV_SESSION_STATE_HOME: Final[str] = f'{HV_STATE_HOME}/sessions'
HV_TOR_STATE_HOME: Final[str] = f'{HV_STATE_HOME}/tor'
VLESS_DNS_ENABLED: Final[bool] = os.environ.get('VLESS_DNS_ENABLED', 'true').lower() == 'true'
HYSTERIA2_DNS_ENABLED: Final[bool] = os.environ.get('HYSTERIA2_DNS_ENABLED', 'true').lower() == 'true'

View file

@ -1,102 +1,56 @@
class CommandNotFoundError(OSError):
def __init__(self, subject):
self.subject = subject
super().__init__(f"Command '{subject}' could not be found.")
class UnknownClientPathError(Exception):
pass
class UnknownClientVersionError(Exception):
pass
class UnknownConnectionTypeError(Exception):
pass
class UnknownTimeZoneError(Exception):
pass
class ConnectionTerminationError(Exception):
pass
class PolicyAssignmentError(Exception):
pass
class PolicyInstatementError(Exception):
pass
class PolicyRevocationError(Exception):
pass
class ProfileDeletionError(Exception):
pass
class ProfileModificationError(Exception):
pass
class ProfileStateConflictError(Exception):
pass
class ProfileActivationError(Exception):
pass
class ProfileDeactivationError(Exception):
pass
class UnsupportedApplicationVersionError(Exception):
pass
class ApplicationAlreadyInstalledError(Exception):
pass
class MissingLocationError(Exception):
pass
class MissingSubscriptionError(Exception):
pass
class InvalidSubscriptionError(Exception):
pass
class InvoiceNotFoundError(Exception):
pass
class InvoiceExpiredError(Exception):
pass
class InvoicePaymentFailedError(Exception):
pass
class ConnectionUnprotectedError(Exception):
pass
class FileIntegrityError(Exception):
pass
class EndpointVerificationError(Exception):
pass
class SingboxNotInstalledException(Exception):
pass
class PaymentRequiredError(Exception):
pass

View file

@ -25,7 +25,7 @@ class ConfigurationController:
configuration = ConfigurationController.get()
if configuration is None or configuration.connection not in ('system', 'tor'):
if configuration is None or configuration.connection not in ('system', 'tor', 'vless', 'hysteria2'):
raise UnknownConnectionTypeError('The preferred connection type could not be determined.')
return configuration.connection

View file

@ -3,7 +3,6 @@ from core.errors.logger import logger
from collections.abc import Callable
from core.Constants import Constants
from core.Errors import InvalidSubscriptionError, MissingSubscriptionError, ConnectionUnprotectedError, ConnectionTerminationError, CommandNotFoundError
from core.controllers.ConfigurationController import ConfigurationController
from core.controllers.ProfileController import ProfileController
from core.controllers.SessionStateController import SessionStateController
@ -16,16 +15,19 @@ from core.services.WebServiceApiService import WebServiceApiService
from essentials.modules.TorModule import TorModule
from essentials.services.ConnectionService import ConnectionService
from pathlib import Path
from core.Errors import InvalidSubscriptionError, MissingSubscriptionError, ConnectionUnprotectedError, ConnectionTerminationError, CommandNotFoundError, PaymentRequiredError
from subprocess import CalledProcessError
from typing import Union, Optional, Any
import os
import random
import re
import shutil
import signal
import subprocess
import sys
import tempfile
import time
from core.models.OperatorProxySession import OperatorProxySession
class ConnectionController:
@ -53,7 +55,10 @@ class ConnectionController:
@staticmethod
def establish_connection(profile: Union[SessionProfile, SystemProfile], ignore: tuple[type[Exception]] = (), connection_observer: Optional[ConnectionObserver] = None):
from core.controllers.ConnectionController import ConnectionController
ConnectionController.verify_singbox_installation()
connection = profile.connection
# needs_proxy_configuration comes from the child connection models
@ -99,6 +104,29 @@ class ConnectionController:
raise MissingSubscriptionError()
if connection.needs_operator_proxy() and not profile.has_operator_proxy_session():
if profile.has_subscription():
if not profile.subscription.has_been_activated():
ProfileController.activate_subscription(profile, connection_observer=connection_observer)
operator_proxy_session = ConnectionController.with_preferred_connection(
profile.subscription.billing_code,
profile.subscription.operator_id,
connection.get_protocol(),
task=WebServiceApiService.post_operator_proxy,
connection_observer=connection_observer
)
if operator_proxy_session is None:
raise PaymentRequiredError()
profile.attach_operator_proxy_session(operator_proxy_session)
else:
raise MissingSubscriptionError()
if profile.is_session_profile():
try:
@ -166,6 +194,24 @@ class ConnectionController:
ConnectionController.establish_wireguard_session_connection(profile, session_directory, port_number)
session_state.network_port_numbers.wireguard.append(port_number)
elif profile.connection.code in ('vless', 'hysteria2'):
if not profile.has_operator_proxy_session():
raise MissingSubscriptionError()
operator_proxy_session = profile.get_operator_proxy_session()
port_number = ConnectionService.get_random_available_port_number()
if profile.connection.code == 'vless':
from core.controllers.encrypted_proxy.VlessController import VlessController
VlessController.enable(operator_proxy_session, port_number)
session_state.network_port_numbers.vless.append(port_number)
elif profile.connection.code == 'hysteria2':
from core.controllers.encrypted_proxy.HysteriaController import HysteriaController
HysteriaController.enable(operator_proxy_session, port_number)
session_state.network_port_numbers.hysteria2.append(port_number)
if profile.connection.masked:
while proxy_port_number is None or proxy_port_number == port_number:
@ -181,8 +227,42 @@ class ConnectionController:
return proxy_port_number or port_number
@staticmethod
def _signal_previous_process() -> None:
system_state = SystemStateController.get()
if system_state is None:
return
pid = system_state.pid
if pid is None or pid == os.getpid():
return
try:
os.kill(pid, signal.SIGUSR1)
time.sleep(0.5)
except (ProcessLookupError, PermissionError):
pass
@staticmethod
def establish_system_connection(profile: SystemProfile, ignore: tuple[type[Exception]] = (), connection_observer: Optional[ConnectionObserver] = None):
ConnectionController._signal_previous_process()
system_state = SystemStateController.get()
if system_state is not None:
try:
ConnectionController.terminate_system_connection(connection_observer=connection_observer)
except ConnectionTerminationError:
pass
if profile.connection.needs_operator_proxy():
ok = ConnectionController.establish_encrypted_proxy_connection(
profile, socks5_port=1080, observer=connection_observer
)
if not ok:
raise ConnectionError('The connection could not be established.')
token = SystemStateController.create(profile.id)
if connection_observer is not None:
connection_observer.notify('connected_token', {'session_token': token})
return
if ConfigurationController.get_endpoint_verification_enabled():
ProfileController.verify_wireguard_endpoint(profile, ignore=ignore)
@ -191,36 +271,63 @@ class ConnectionController:
ConnectionController.__establish_system_connection(profile, connection_observer)
except ConnectionError:
try:
ConnectionController.terminate_system_connection()
except ConnectionTerminationError:
pass
raise ConnectionError('The connection could not be established.')
except CalledProcessError:
try:
ConnectionController.terminate_system_connection()
except ConnectionTerminationError:
pass
try:
ConnectionController.__establish_system_connection(profile, connection_observer)
except (ConnectionError, CalledProcessError):
try:
ConnectionController.terminate_system_connection()
except ConnectionTerminationError:
pass
raise ConnectionError('The connection could not be established.')
ConnectionController.terminate_tor_connection()
time.sleep(1.0)
@staticmethod
def establish_encrypted_proxy_connection(profile, socks5_port: int = 1080, observer=None) -> bool:
"""Shared encrypted proxy connection logic for CLI and GUI."""
import socket
operator_proxy_session = profile.get_operator_proxy_session()
protocol = profile.connection.code
if protocol == 'vless':
from core.controllers.encrypted_proxy.VlessController import VlessController
from core.services.encrypted_proxy.vless_service import parse_vless_link
server_ip = operator_proxy_session.server_ip
if server_ip is None:
vless = parse_vless_link(operator_proxy_session.links[0])
server_ip = socket.gethostbyname(vless['host'])
return VlessController(socks5_port).enable(
operator_proxy_session.links[0],
operator_proxy_session.username,
server_ip,
observer
)
elif protocol == 'hysteria2':
from core.controllers.encrypted_proxy.HysteriaController import HysteriaController
server_ip = operator_proxy_session.server_ip
if server_ip is None:
server_ip = socket.gethostbyname(operator_proxy_session.operator_hysteria2_host)
return HysteriaController(socks5_port).enable(
operator_proxy_session.username,
operator_proxy_session.password,
operator_proxy_session.operator_hysteria2_host,
server_ip,
observer
)
return False
@staticmethod
def establish_tor_connection(connection_observer: Optional[ConnectionObserver] = None):
@ -309,22 +416,39 @@ class ConnectionController:
return subprocess.Popen(('proxychains4', '-f', proxychains_configuration_file_path, 'microsocks', '-p', str(proxy_port_number)), stdout=subprocess.DEVNULL, stderr=subprocess.STDOUT)
@staticmethod
def terminate_system_connection():
def terminate_system_connection(connection_observer: Optional[ConnectionObserver] = None):
system_state = SystemStateController.get()
if system_state is not None:
profile = ProfileController.get(system_state.profile_id)
if profile is not None and profile.connection.needs_operator_proxy():
protocol = profile.connection.get_protocol()
if protocol == 'vless':
from core.controllers.encrypted_proxy.VlessController import VlessController
VlessController(1080).disable(connection_observer)
elif protocol == 'hysteria2':
from core.controllers.encrypted_proxy.HysteriaController import HysteriaController
HysteriaController(1080).disable(connection_observer)
SystemState.dissolve()
return
if shutil.which('nmcli') is None:
raise CommandNotFoundError('nmcli')
if SystemStateController.exists():
process = subprocess.Popen(('nmcli', 'connection', 'delete', 'wg'), stdout=subprocess.DEVNULL, stderr=subprocess.STDOUT)
completed_successfully = not bool(os.waitpid(process.pid, 0)[1] >> 8)
if completed_successfully or not ConnectionController.system_uses_wireguard_interface():
subprocess.run(('nmcli', 'connection', 'delete', 'hv-ipv6-sink'), stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
ConnectionController.terminate_tor_connection()
try:
from core.utils.encrypted_proxy import killswitch
killswitch.disarm()
except Exception:
pass
SystemState.dissolve()
if connection_observer is not None:
connection_observer.notify('disconnected', {})
else:
raise ConnectionTerminationError('The connection could not be terminated.')
@ -430,7 +554,18 @@ class ConnectionController:
except CalledProcessError:
raise ConnectionError('The connection could not be established.')
SystemStateController.create(profile.id)
# Arm killswitch for WireGuard
try:
wg_server_ip = ConnectionController.__extract_wireguard_endpoint(profile)
if wg_server_ip:
from core.utils.encrypted_proxy import killswitch
killswitch.arm(wg_server_ip, 'wg')
except Exception:
pass
token = SystemStateController.create(profile.id)
if connection_observer is not None:
connection_observer.notify('connected_token', {'session_token': token})
try:
ConnectionController.await_connection(connection_observer=connection_observer)
@ -438,6 +573,23 @@ class ConnectionController:
except ConnectionError:
raise ConnectionError('The connection could not be established.')
@staticmethod
def __extract_wireguard_endpoint(profile):
import re, socket
try:
config = profile.get_wireguard_configuration()
if config is None:
return None
match = re.search(r'Endpoint\s*=\s*([^\s:]+):\d+', config)
if not match:
return None
host = match.group(1)
if re.match(r'^(\d{1,3}\.){3}\d{1,3}$', host):
return host
return socket.gethostbyname(host)
except Exception:
return None
@staticmethod
def __with_tor_connection(*args, task: Callable[..., Any], connection_observer: Optional[ConnectionObserver] = None, **kwargs):
@ -503,3 +655,25 @@ class ConnectionController:
return True
return False
@staticmethod
def verify_singbox_installation():
import subprocess
import os
from core.Errors import SingboxNotInstalledException
singbox_binary = '/usr/bin/sing-box'
singbox_wrapper = '/usr/local/bin/hydraveil-singbox'
if not os.path.exists(singbox_binary):
raise SingboxNotInstalledException(f'sing-box binary not found at {singbox_binary}')
if not os.path.exists(singbox_wrapper):
raise SingboxNotInstalledException(f'sing-box wrapper not found at {singbox_wrapper}')
try:
subprocess.run([singbox_binary, 'version'], capture_output=True, timeout=5, check=True)
except subprocess.CalledProcessError:
raise SingboxNotInstalledException('sing-box binary check failed')
except subprocess.TimeoutExpired:
raise SingboxNotInstalledException('sing-box verification timeout')

View file

@ -1,4 +1,4 @@
from core.Errors import InvalidSubscriptionError, MissingSubscriptionError, ConnectionTerminationError, ProfileActivationError, ProfileDeactivationError, MissingLocationError, ConnectionUnprotectedError, EndpointVerificationError, ProfileStateConflictError
from core.Errors import InvalidSubscriptionError, MissingSubscriptionError, ConnectionTerminationError, ProfileActivationError, ProfileDeactivationError, MissingLocationError, ConnectionUnprotectedError, EndpointVerificationError, ProfileStateConflictError, PaymentRequiredError
from core.controllers.ApplicationController import ApplicationController
from core.controllers.ApplicationVersionController import ApplicationVersionController
from core.controllers.SessionStateController import SessionStateController
@ -77,6 +77,14 @@ class ProfileController:
if profile.is_system_profile():
if profile.connection.code in ('vless', 'hysteria2'):
ProfileController._enable_encrypted_proxy(profile, connection_observer=connection_observer)
if profile_observer is not None:
profile_observer.notify('enabled', profile)
return
# WireGuard and other system profiles — legacy flow
from core.controllers.ConnectionController import ConnectionController
try:
ConnectionController.establish_connection(profile, ignore=ignore, connection_observer=connection_observer)
except ConnectionError:
@ -86,51 +94,73 @@ class ProfileController:
profile_observer.notify('enabled', profile)
@staticmethod
def disable(profile: Union[SessionProfile, SystemProfile], explicitly: bool = True, ignore: tuple[type[Exception]] = (), profile_observer: ProfileObserver = None):
def _enable_encrypted_proxy(profile, connection_observer=None) -> None:
from core.controllers.ConnectionController import ConnectionController
from core.controllers.SystemStateController import SystemStateController
from core.services.WebServiceApiService import WebServiceApiService
if not profile.has_operator_proxy_session():
ProfileController.activate_subscription(profile, connection_observer=connection_observer)
operator_proxy_session = ConnectionController.with_preferred_connection(
profile.subscription.billing_code,
profile.subscription.operator_id,
profile.connection.code,
task=WebServiceApiService.post_operator_proxy,
connection_observer=connection_observer
)
if operator_proxy_session is None:
raise ProfileActivationError('Could not obtain operator proxy session.')
profile.attach_operator_proxy_session(operator_proxy_session)
ConnectionController._signal_previous_process()
ok = ConnectionController.establish_encrypted_proxy_connection(
profile, socks5_port=1080, observer=connection_observer
)
if not ok:
raise ProfileActivationError('The profile could not be enabled.')
token = SystemStateController.create(profile.id)
if connection_observer is not None:
connection_observer.notify('connected_token', {'session_token': token})
@staticmethod
def run_monitor(profile: Union[SessionProfile, SystemProfile], monitor) -> None:
if profile.is_system_profile() and profile.connection.needs_operator_proxy():
monitor.run_blocking()
@staticmethod
def disable(profile: Union[SessionProfile, SystemProfile], explicitly: bool = True, ignore: tuple[type[Exception]] = (), profile_observer: ProfileObserver = None, connection_observer: ConnectionObserver = None):
from core.controllers.ConnectionController import ConnectionController
if profile.is_session_profile():
if SessionStateController.exists(profile.id):
session_state = SessionStateController.get(profile.id)
if session_state is not None:
for port_number in session_state.network_port_numbers.tor:
ConnectionController.terminate_tor_session_connection(port_number)
session_state.dissolve(session_state.id)
if profile.is_system_profile():
subjects = ProfileController.get_all().values()
for subject in subjects:
if subject.is_session_profile():
if subject.connection.is_unprotected() and ProfileController.is_enabled(subject) and not ConnectionUnprotectedError in ignore:
raise ConnectionUnprotectedError('Disabling this system connection would leave one or more sessions exposed.')
if SystemStateController.exists():
system_state = SystemStateController.get()
if profile.id != system_state.profile_id:
raise ProfileDeactivationError('The profile could not be disabled.')
try:
ConnectionController.terminate_system_connection()
ConnectionController.terminate_system_connection(connection_observer=connection_observer)
except ConnectionTerminationError:
raise ProfileDeactivationError('The profile could not be disabled.')
if profile_observer is not None:
profile_observer.notify('disabled', profile, dict(
explicitly=explicitly,
))
profile_observer.notify('disabled', profile, dict(explicitly=explicitly))
time.sleep(1.0)
@ -156,13 +186,29 @@ class ProfileController:
if profile.has_subscription():
# Operator profiles (vless/hysteria2): check if already activated locally first.
# If not, poll the server — the payment may have been processed since last attempt.
if profile.subscription.operator_id is not None:
if profile.subscription.has_been_activated():
profile.save()
return
subscription = ConnectionController.with_preferred_connection(
profile.subscription.billing_code,
task=WebServiceApiService.get_subscription,
connection_observer=connection_observer
)
if subscription is not None:
profile.subscription = subscription
profile.save()
return
else:
raise PaymentRequiredError()
subscription = ConnectionController.with_preferred_connection(profile.subscription.billing_code, task=WebServiceApiService.get_subscription, connection_observer=connection_observer)
if subscription is not None:
profile.subscription = subscription
profile.save()
else:
raise InvalidSubscriptionError()
@ -310,4 +356,4 @@ class ProfileController:
return dict(
private=base64.b64encode(private_key).decode(),
public=base64.b64encode(public_key).decode()
)
)

View file

@ -5,7 +5,6 @@ from core.observers.ConnectionObserver import ConnectionObserver
from core.services.WebServiceApiService import WebServiceApiService
from typing import Union
class SubscriptionController:
@staticmethod
@ -18,4 +17,51 @@ class SubscriptionController:
def create(subscription_plan: SubscriptionPlan, profile: Union[SessionProfile, SystemProfile], connection_observer: ConnectionObserver = None):
from core.controllers.ConnectionController import ConnectionController
return ConnectionController.with_preferred_connection(subscription_plan.id, profile.location.id, task=WebServiceApiService.post_subscription, connection_observer=connection_observer)
if profile.location:
return ConnectionController.with_preferred_connection(
subscription_plan.id,
profile.location.id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)
else:
return ConnectionController.with_preferred_connection(
subscription_plan.id,
operator_id=profile.connection.operator_id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)
from core.controllers.ConnectionController import ConnectionController
if profile.location:
return ConnectionController.with_preferred_connection(
subscription_plan.id,
profile.location.id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)
else:
return ConnectionController.with_preferred_connection(
subscription_plan.id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)
from core.controllers.ConnectionController import ConnectionController
if profile.location:
# Para WireGuard (con location)
return ConnectionController.with_preferred_connection(
subscription_plan.id,
profile.location.id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)
else:
return ConnectionController.with_preferred_connection(
subscription_plan.id,
task=WebServiceApiService.post_subscription,
connection_observer=connection_observer
)

View file

@ -1,3 +1,5 @@
import os
import uuid
from core.models.system.SystemState import SystemState
@ -13,7 +15,10 @@ class SystemStateController:
@staticmethod
def create(profile_id):
return SystemState(profile_id).save()
token = str(uuid.uuid4())
pid = os.getpid()
SystemState(profile_id, token, pid).save()
return token
@staticmethod
def update_or_create(system_state):
@ -21,4 +26,4 @@ class SystemStateController:
@staticmethod
def dissolve():
return SystemState.dissolve()
return SystemState.dissolve()

View file

@ -0,0 +1,17 @@
from pathlib import Path
from core.services.encrypted_proxy.disable_service import disable_proxy
class DisableController:
def __init__(self, tmp_dir: Path, wrapper: str, unit: str):
self.tmp_dir = tmp_dir
self.wrapper = wrapper
self.unit = unit
def disable(self, observer=None) -> bool:
return disable_proxy(
tmp_dir=self.tmp_dir,
wrapper=self.wrapper,
unit=self.unit,
observer=observer,
)

View file

@ -0,0 +1,31 @@
from core.services.encrypted_proxy.hysteria_service import enable_hysteria, disable_hysteria
class HysteriaController:
def __init__(self, socks5_port: int):
self.socks5_port = socks5_port
def enable(self, username: str, password: str,
server_host: str, server_ip: str, observer=None) -> bool:
if not username or not isinstance(username, str):
if observer:
observer.notify("error", "Invalid username")
return False
if not password or not server_host:
if observer:
observer.notify("error", "Missing password or server_host")
return False
if not server_ip:
if observer:
observer.notify("error", "Missing server_ip")
return False
return enable_hysteria(
username=username,
password=password,
server_host=server_host,
server_ip=server_ip,
socks5_port=self.socks5_port,
observer=observer,
)
def disable(self, observer=None) -> bool:
return disable_hysteria(observer=observer)

View file

@ -0,0 +1,29 @@
from core.services.encrypted_proxy.vless_service import enable_vless, disable_vless
class VlessController:
def __init__(self, socks5_port: int):
self.socks5_port = socks5_port
def enable(self, vless_link: str, username: str, server_ip: str, observer=None) -> bool:
if not username or not isinstance(username, str):
if observer:
observer.notify("error", "Invalid username")
return False
if not vless_link or not vless_link.startswith("vless://"):
if observer:
observer.notify("error", "Invalid vless link")
return False
if not server_ip:
if observer:
observer.notify("error", "Missing server_ip")
return False
return enable_vless(
vless_link=vless_link,
username=username,
server_ip=server_ip,
socks5_port=self.socks5_port,
observer=observer,
)
def disable(self, observer=None) -> bool:
return disable_vless(observer=observer)

View file

@ -1,19 +1,15 @@
from dataclasses import dataclass
from dataclasses_json import dataclass_json
@dataclass_json
@dataclass
class BaseConnection:
code: str
# called by: connection controller for basic type checks on wireguard code
def needs_wireguard_configuration(self):
return self.code == 'wireguard'
def is_session_connection(self):
return type(self).__name__ == 'SessionConnection'
def is_system_connection(self):
return type(self).__name__ == 'SystemConnection'
def needs_operator_proxy(self):
return False

View file

@ -65,11 +65,10 @@ class BaseProfile(ABC):
return type(self).__name__ == 'SystemProfile'
def save(self: Self):
config_dict = self.to_dict() # Get dict, not JSON string
location_dict = self.location.to_dict() # this is from SQLAlchemy, and not JSON-models, that's why it's separate.
config_dict = self.to_dict()
location_dict = self.location.to_dict() if self.location else None
if self.location:
config_dict["location"] = location_dict
config_dict["location"] = location_dict
config_file_contents = json.dumps(config_dict, indent=4) + '\n'

View file

@ -0,0 +1,22 @@
from dataclasses import dataclass
from dataclasses_json import dataclass_json
from typing import Optional
@dataclass_json
@dataclass
class OperatorProxySession:
id: int
type: str
username: Optional[str]
password: Optional[str]
links: Optional[list]
subscription_url: Optional[str]
operator_id: int
operator_name: str
operator_domain: Optional[str] = None
operator_hysteria2_host: Optional[str] = None
operator_vless_host: Optional[str] = None
server_ip: Optional[str] = None # pre-resolved IP — avoids DNS leak at connect time
location_country_code: Optional[str] = None
location_city_code: Optional[str] = None

View file

@ -10,6 +10,13 @@ import dataclasses_json
@dataclass
class Subscription:
billing_code: str
operator_id: Optional[int] = field(
default=None,
metadata=config(
undefined=dataclasses_json.Undefined.EXCLUDE,
exclude=lambda value: value is None
)
)
expires_at: Optional[datetime] = field(
default=None,
metadata=config(
@ -37,4 +44,4 @@ class Subscription:
@staticmethod
def _iso_format(datetime_instance: datetime):
return datetime.isoformat(datetime_instance).replace('+00:00', 'Z')
return datetime.isoformat(datetime_instance).replace('+00:00', 'Z')

View file

@ -53,6 +53,9 @@ class SubscriptionPlan(Model):
if connection.code == 'wireguard':
features_wireguard = True
if connection.code in ('operator', 'vless', 'hysteria2'):
features_proxy = True
Model._create_table_if_not_exists(table_name=_table_name, table_definition=_table_definition)
return Model._query_one('SELECT * FROM subscription_plans WHERE features_proxy = ? AND features_wireguard = ? AND duration = ? LIMIT 1', SubscriptionPlan.factory, [features_proxy, features_wireguard, duration])
@ -76,6 +79,9 @@ class SubscriptionPlan(Model):
if connection.code == 'wireguard':
features_wireguard = True
if connection.code in ('operator', 'vless', 'hysteria2'):
features_proxy = True
return Model._query_all('SELECT * FROM subscription_plans WHERE features_proxy = ? AND features_wireguard = ?', SubscriptionPlan.factory, [features_proxy, features_wireguard])
@staticmethod
@ -94,4 +100,4 @@ class SubscriptionPlan(Model):
@staticmethod
def tuple_factory(subscription_plan):
return subscription_plan.id, subscription_plan.code, subscription_plan.wireguard_session_limit, subscription_plan.duration, subscription_plan.price, subscription_plan.features_proxy, subscription_plan.features_wireguard
return subscription_plan.id, subscription_plan.code, subscription_plan.wireguard_session_limit, subscription_plan.duration, subscription_plan.price, subscription_plan.features_proxy, subscription_plan.features_wireguard

View file

@ -6,11 +6,13 @@ class NetworkPortNumbers:
proxy: list[int] = field(default_factory=list)
wireguard: list[int] = field(default_factory=list)
tor: list[int] = field(default_factory=list)
vless: list[int] = field(default_factory=list)
hysteria2: list[int] = field(default_factory=list)
@property
def all(self):
return self.proxy + self.wireguard + self.tor
return self.proxy + self.wireguard + self.tor + self.vless + self.hysteria2
@property
def isolated(self):
return self.proxy + self.wireguard
return self.proxy + self.wireguard + self.vless + self.hysteria2

View file

@ -1,14 +1,12 @@
from core.models.BaseConnection import BaseConnection
from dataclasses import dataclass
@dataclass
class SessionConnection(BaseConnection):
masked: bool = False
def __post_init__(self):
if self.code not in ('system', 'tor', 'wireguard'):
if self.code not in ('system', 'tor', 'wireguard', 'vless', 'hysteria2'):
raise ValueError('Invalid connection code.')
# called by connection controller
@ -18,3 +16,9 @@ class SessionConnection(BaseConnection):
def needs_proxy_configuration(self):
return self.masked is True
def needs_operator_proxy(self):
return self.code in ('vless', 'hysteria2')
def get_protocol(self):
return self.code if self.needs_operator_proxy() else None

View file

@ -22,35 +22,24 @@ class SessionProfile(BaseProfile):
return self.connection is not None
def save(self):
if 'application_version' in self._get_dirty_keys():
persistent_state_path = f'{self.get_data_path()}/persistent-state'
if os.path.isdir(persistent_state_path):
shutil.rmtree(persistent_state_path, ignore_errors=True)
if 'location' in self._get_dirty_keys():
self.__delete_proxy_configuration()
self.__delete_wireguard_configuration()
super().save()
def attach_proxy_configuration(self, proxy_configuration):
proxy_configuration_file_contents = f'{proxy_configuration.to_json(indent=4)}\n'
os.makedirs(Constants.HV_CONFIG_HOME, exist_ok=True)
proxy_configuration_file_path = self.get_proxy_configuration_path()
with open(proxy_configuration_file_path, 'w') as proxy_configuration_file:
proxy_configuration_file.write(proxy_configuration_file_contents)
def attach_wireguard_configuration(self, wireguard_configuration):
wireguard_configuration_file_path = self.get_wireguard_configuration_path()
with open(wireguard_configuration_file_path, 'w') as wireguard_configuration_file:
wireguard_configuration_file.write(wireguard_configuration)
@ -61,19 +50,15 @@ class SessionProfile(BaseProfile):
return f'{self.get_config_path()}/wg.conf'
def get_proxy_configuration(self):
try:
config_file_contents = open(self.get_proxy_configuration_path(), 'r').read()
except FileNotFoundError:
return None
try:
proxy_configuration = json.loads(config_file_contents)
except ValueError:
return None
proxy_configuration = ProxyConfiguration.from_dict(proxy_configuration)
return proxy_configuration
def has_proxy_configuration(self):
@ -83,36 +68,26 @@ class SessionProfile(BaseProfile):
return os.path.isfile(f'{self.get_config_path()}/wg.conf')
def address_security_incident(self):
super().address_security_incident()
self.__delete_wireguard_configuration()
def determine_timezone(self):
time_zone = None
if self.has_connection():
if self.connection.needs_proxy_configuration():
if self.has_proxy_configuration():
time_zone = self.get_proxy_configuration().time_zone
elif self.connection.needs_wireguard_configuration():
if self.has_wireguard_configuration():
time_zone = self.get_wireguard_configuration_metadata('TZ')
if time_zone is None and self.has_location():
time_zone = self.location.time_zone
if time_zone is None:
raise UnknownTimeZoneError('The preferred time zone could not be determined.')
return time_zone
def __delete_proxy_configuration(self):
Path(self.get_proxy_configuration_path()).unlink(missing_ok=True)
def __delete_wireguard_configuration(self):
Path(self.get_wireguard_configuration_path()).unlink(missing_ok=True)
Path(self.get_wireguard_configuration_path()).unlink(missing_ok=True)

View file

@ -1,16 +1,22 @@
from core.models.BaseConnection import BaseConnection
from dataclasses import dataclass
from dataclasses import dataclass, field
from typing import Optional
@dataclass
class SystemConnection (BaseConnection):
class SystemConnection(BaseConnection):
operator_id: Optional[int] = field(default=None)
protocol: Optional[str] = field(default=None)
def __post_init__(self):
if self.code not in ('vless', 'hysteria2', 'wireguard'):
if self.code not in ('wireguard', 'operator', 'vless', 'hysteria2'):
raise ValueError('Invalid connection code.')
@staticmethod
def needs_proxy_configuration():
return False
def needs_operator_proxy(self):
return self.code in ('operator', 'vless', 'hysteria2')
def get_protocol(self):
return self.protocol if self.protocol else self.code

View file

@ -4,11 +4,11 @@ from core.models.BaseProfile import BaseProfile
from core.models.system.SystemConnection import SystemConnection
from dataclasses import dataclass
from typing import Optional
import json
import os
import shutil
import subprocess
@dataclass
class SystemProfile(BaseProfile):
connection: Optional[SystemConnection]
@ -19,33 +19,23 @@ class SystemProfile(BaseProfile):
return filepath
def save(self):
if 'location' in self._get_dirty_keys():
self.__delete_wireguard_configuration()
super().save()
def attach_wireguard_configuration(self, wireguard_configuration):
if shutil.which('pkexec') is None:
raise CommandNotFoundError('pkexec')
wireguard_configuration_file_backup_path = f'{self.get_config_path()}/wg.conf.bak'
with open(wireguard_configuration_file_backup_path, 'w') as wireguard_configuration_file:
wireguard_configuration_file.write(wireguard_configuration)
wireguard_configuration_is_attached = False
failed_attempt_count = 0
while not wireguard_configuration_is_attached and failed_attempt_count < 3:
process = subprocess.Popen(('pkexec', 'install', '-D', wireguard_configuration_file_backup_path, self.get_wireguard_configuration_path(), '-o', 'root', '-m', '744'))
wireguard_configuration_is_attached = not bool(os.waitpid(process.pid, 0)[1] >> 8)
if not wireguard_configuration_is_attached:
failed_attempt_count += 1
if not wireguard_configuration_is_attached:
raise ProfileModificationError('The WireGuard configuration could not be attached.')
@ -61,7 +51,6 @@ class SystemProfile(BaseProfile):
return False
def address_security_incident(self):
super().address_security_incident()
self.__delete_wireguard_configuration()
@ -70,7 +59,6 @@ class SystemProfile(BaseProfile):
self.__delete_wireguard_configuration()
except ProfileModificationError:
raise ProfileDeletionError('The WireGuard configuration could not be deleted.')
if shutil.which('pkexec') is None:
raise CommandNotFoundError('pkexec')
@ -84,19 +72,59 @@ class SystemProfile(BaseProfile):
super().delete()
def attach_operator_proxy_session(self, operator_proxy_session):
from core.models.OperatorProxySession import OperatorProxySession
operator_proxy_session_file_contents = f'{operator_proxy_session.to_json(indent=4)}\n'
os.makedirs(self.get_config_path(), exist_ok=True)
operator_proxy_session_file_path = self.get_operator_proxy_session_path()
with open(operator_proxy_session_file_path, 'w') as operator_proxy_session_file:
operator_proxy_session_file.write(operator_proxy_session_file_contents)
if operator_proxy_session.location_country_code and operator_proxy_session.location_city_code:
try:
from core.models.orm_models.Location import Location
from core.models.manage.session_management import get_session
from sqlalchemy import select
session = get_session()
loc = session.execute(
select(Location).where(
(Location.country_code == operator_proxy_session.location_country_code) &
(Location.code == operator_proxy_session.location_city_code)
)
).scalar_one_or_none()
if loc:
self.location = loc
self.save()
except Exception:
pass
def get_operator_proxy_session_path(self):
return f'{self.get_config_path()}/operator_proxy_session.json'
def get_operator_proxy_session(self):
try:
config_file_contents = open(self.get_operator_proxy_session_path(), 'r').read()
except FileNotFoundError:
return None
try:
data = json.loads(config_file_contents)
except ValueError:
return None
from core.models.OperatorProxySession import OperatorProxySession
return OperatorProxySession.from_dict(data)
def has_operator_proxy_session(self):
return os.path.isfile(self.get_operator_proxy_session_path())
def __delete_wireguard_configuration(self):
if self.has_wireguard_configuration():
if shutil.which('pkexec') is None:
raise CommandNotFoundError('pkexec')
try:
process = subprocess.run(('pkexec', 'rm', '-rf', self.get_wireguard_configuration_path()), check=True)
completed_successfully = not bool(os.waitpid(process.pid, 0)[1] >> 8)
except subprocess.CalledProcessError as e:
completed_successfully = True
completed_successfully = True
except:
completed_successfully = True

View file

@ -1,8 +1,8 @@
from core.Constants import Constants
from dataclasses import dataclass
from dataclasses import dataclass, field
from dataclasses_json import dataclass_json
from pathlib import Path
from typing import Self
from typing import Optional, Self
import json
import os
import pathlib
@ -11,6 +11,8 @@ import pathlib
@dataclass
class SystemState:
profile_id: int
session_token: Optional[str] = field(default=None)
pid: Optional[int] = field(default=None)
def save(self: Self):
@ -31,7 +33,6 @@ class SystemState:
system_state_file_contents = open(f'{SystemState.__get_state_path()}/system.json', 'r').read()
system_state_dict = json.loads(system_state_file_contents)
# noinspection PyUnresolvedReferences
return SystemState.from_dict(system_state_dict)
except (FileNotFoundError, ValueError, KeyError):
@ -53,4 +54,4 @@ class SystemState:
@staticmethod
def __get_state_path():
return Constants.HV_STATE_HOME
return Constants.HV_STATE_HOME

View file

@ -1,5 +1,9 @@
from essentials.observers.ConnectionObserver import ConnectionObserver as BaseConnectionObserver
class ConnectionObserver(BaseConnectionObserver):
pass
def __init__(self):
super().__init__()
self.on_connected = []
self.on_disconnected = []
self.on_error = []

View file

@ -0,0 +1,8 @@
from core.observers.BaseObserver import BaseObserver
class EncryptedProxyObserver(BaseObserver):
def __init__(self):
self.on_connected = []
self.on_disconnected = []
self.on_error = []

View file

@ -1,7 +1,8 @@
from core.Constants import Constants
from core.models.ClientVersion import ClientVersion
# from core.models.Location import Location
# from core.models.Operator import Operator
# from core.models.Location import Location # migrated to ORM
# from core.models.Operator import Operator # migrated to ORM
from core.models.OperatorProxySession import OperatorProxySession
from core.models.Subscription import Subscription
from core.models.SubscriptionPlan import SubscriptionPlan
from core.models.invoice.Invoice import Invoice
@ -18,12 +19,10 @@ class WebServiceApiService:
@staticmethod
def get_applications(proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/platforms/linux-x86_64/applications', None, proxies)
applications = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for application in response.json()['data']:
applications.append(Application(application['code'], application['name'], application['id']))
@ -32,12 +31,10 @@ class WebServiceApiService:
@staticmethod
def get_application_versions(code: str, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get(f'/platforms/linux-x86_64/applications/{code}/application-versions', None, proxies)
application_versions = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for application_version in response.json()['data']:
application_versions.append(ApplicationVersion(code, application_version['version_number'], application_version['format_revision'], application_version['id'], application_version['download_path'], application_version['released_at'], application_version['file_hash']))
@ -46,12 +43,10 @@ class WebServiceApiService:
@staticmethod
def get_client_versions(proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/platforms/linux-x86_64/appimage/client-versions', None, proxies)
client_versions = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for client_version in response.json()['data']:
client_versions.append(ClientVersion(client_version['version_number'], client_version['released_at'], client_version['id'], client_version['download_path']))
@ -60,26 +55,22 @@ class WebServiceApiService:
@staticmethod
def get_operators(proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/operators', None, proxies)
operators = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for operator in response.json()['data']:
operators.append(Operator(operator['id'], operator['name'], operator['public_key'], operator['nostr_public_key'], operator['nostr_profile_reference'], operator['nostr_attestation']['event_reference']))
operators.append(Operator(operator['id'], operator['name'], operator['type'], operator['public_key'], operator['nostr_public_key'], operator['nostr_profile_reference'], operator['nostr_attestation']['event_reference']))
return operators
@staticmethod
def get_locations(proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/locations', None, proxies)
locations = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for location in response.json()['data']:
locations.append(Location(location['country']['code'], location['code'], location['id'], location['country']['name'], location['name'], location['time_zone']['code'], location['operator_id'], location['provider']['name'], location['is_proxy_capable'], location['is_wireguard_capable']))
@ -88,60 +79,57 @@ class WebServiceApiService:
@staticmethod
def get_subscription_plans(proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/subscription-plans', None, proxies)
subscription_plans = []
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
for subscription_plan in response.json()['data']:
subscription_plans.append(SubscriptionPlan(subscription_plan['id'], subscription_plan['code'], subscription_plan['wireguard_session_limit'], subscription_plan['duration'], subscription_plan['price'], subscription_plan['features_proxy'], subscription_plan['features_wireguard']))
return subscription_plans
@staticmethod
def post_subscription(subscription_plan_id, location_id, proxies: Optional[dict] = None):
def post_subscription(subscription_plan_id, location_id=None, operator_id=None, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__post('/subscriptions', None, {
'subscription_plan_id': subscription_plan_id,
'location_id': location_id
}, proxies)
if response.status_code == status_codes.CREATED:
return Subscription(response.headers['X-Billing-Code'])
body = {'subscription_plan_id': subscription_plan_id}
if operator_id is not None:
body['operator_id'] = operator_id
else:
return None
body['location_id'] = location_id
response = WebServiceApiService.__post('/subscriptions', None, body, proxies)
if 200 <= response.status_code < 300:
return Subscription(response.headers['X-Billing-Code'], operator_id=operator_id)
return None
@staticmethod
def get_subscription(billing_code: str, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
billing_code = billing_code.replace('-', '').upper()
billing_code_fragments = re.findall('....?', billing_code)
billing_code = '-'.join(billing_code_fragments)
response = WebServiceApiService.__get('/subscriptions/current', billing_code, proxies)
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
subscription = response.json()['data']
return Subscription(billing_code, Subscription.from_iso_format(subscription['expires_at']))
return Subscription(
billing_code,
operator_id=subscription.get('operator_id'),
expires_at=Subscription.from_iso_format(subscription['expires_at'])
)
else:
return None
return None
@staticmethod
def get_invoice(billing_code: str, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/invoices/current', billing_code, proxies)
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
response_data = response.json()['data']
@ -157,37 +145,59 @@ class WebServiceApiService:
return Invoice(billing_code, invoice['status'], invoice['expires_at'], tuple[PaymentMethod](payment_methods))
else:
return None
return None
@staticmethod
def get_proxy_configuration(billing_code: str, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__get('/proxy-configurations/current', billing_code, proxies)
if response.status_code == status_codes.OK:
if 200 <= response.status_code < 300:
proxy_configuration = response.json()['data']
return ProxyConfiguration(proxy_configuration['ip_address'], proxy_configuration['port'], proxy_configuration['username'], proxy_configuration['password'], proxy_configuration['location']['time_zone']['code'])
else:
return None
return None
@staticmethod
def post_operator_proxy(billing_code: str, operator_id: int, protocol: str, proxies: Optional[dict] = None):
response = WebServiceApiService.__post('/subscriptions/current/operator-proxies', billing_code, {
'operator_id': operator_id,
'protocol': protocol,
}, proxies)
if 200 <= response.status_code < 300:
data = response.json()['data']
return OperatorProxySession(
data['id'],
data['type'],
data['username'],
data.get('password'),
data.get('links'),
data.get('subscription_url'),
data['operator']['id'],
data['operator']['name'],
data['operator'].get('domain'),
data['operator'].get('hysteria2_host'),
data['operator'].get('vless_host'),
data.get('server_ip'),
data.get('location_country_code'),
data.get('location_city_code'),
)
return None
@staticmethod
def post_wireguard_session(country_code: str, location_code: str, billing_code: str, public_key: str, proxies: Optional[dict] = None):
from requests.status_codes import codes as status_codes
response = WebServiceApiService.__post(f'/countries/{country_code}/locations/{location_code}/wireguard-sessions', billing_code, {
'public_key': public_key,
}, proxies)
if response.status_code == status_codes.CREATED:
if 200 <= response.status_code < 300:
return response.text
else:
return None
return None
@staticmethod
def __get(path, billing_code: Optional[str] = None, proxies: Optional[dict] = None):
@ -199,7 +209,7 @@ class WebServiceApiService:
else:
headers = None
return requests.get(Constants.SP_API_BASE_URL + path, headers=headers, proxies=proxies)
return requests.get(Constants.SP_API_BASE_URL + path, headers=headers, proxies=proxies, timeout=30)
@staticmethod
def __post(path, billing_code: Optional[str] = None, body: Optional[dict] = None, proxies: Optional[dict] = None):
@ -211,7 +221,7 @@ class WebServiceApiService:
else:
headers = None
return requests.post(Constants.SP_API_BASE_URL + path, headers=headers, json=body, proxies=proxies)
return requests.post(Constants.SP_API_BASE_URL + path, headers=headers, json=body, proxies=proxies, timeout=30)
@staticmethod
def get_cached_sync(proxies: Optional[dict] = None):

View file

@ -0,0 +1,26 @@
import subprocess
import time
from core.utils.encrypted_proxy.singbox import SingboxRunner
from core.utils.encrypted_proxy.dns import revert_dns_on_tun
from core.utils.encrypted_proxy import killswitch
def get_public_ip(timeout: int = 8) -> str:
endpoints = [
"https://api.ipify.org",
"https://ifconfig.me/ip",
"https://icanhazip.com",
]
for url in endpoints:
try:
r = subprocess.run(
["curl", "-s", "--max-time", str(timeout), url],
capture_output=True,
timeout=timeout + 2,
)
ip = r.stdout.decode().strip()
if ip and "." in ip and not ip.startswith("unknown"):
return ip
except Exception:
continue
return "unknown"

View file

@ -0,0 +1,174 @@
import time
from pathlib import Path
from core.Constants import Constants
from core.utils.encrypted_proxy.singbox import SingboxRunner
from core.utils.encrypted_proxy.dns import wait_for_tun, set_dns_on_tun, revert_dns_on_tun
from core.utils.encrypted_proxy import killswitch
def build_hysteria_config(username: str, password: str,
server_host: str, socks5_port: int,
server_ip: str) -> dict:
return {
"dns": {
"servers": [{"tag": "tunnel-dns", "type": "udp", "server": "9.9.9.9"}],
"final": "tunnel-dns",
"strategy": "ipv4_only",
"independent_cache": True,
},
"inbounds": [
{
"type": "tun",
"tag": "tun-in",
"interface_name": Constants.SINGBOX_TUN_IF,
"address": [Constants.SINGBOX_INTERNAL_ADDR],
"mtu": 9000,
"auto_route": True,
"stack": "gvisor",
},
{
"type": "socks",
"tag": "socks-in",
"listen": "127.0.0.1",
"listen_port": socks5_port,
},
],
"outbounds": [
{"type": "direct", "tag": "direct"},
{"type": "block", "tag": "block"},
{
"type": "hysteria2",
"tag": "proxy",
"server": server_ip,
"server_port": 443,
"password": f"{username}:{password}",
"tls": {
"enabled": True,
"server_name": server_host,
"insecure": False,
},
},
],
"route": {
"rules": [
{"protocol": "dns", "action": "hijack-dns"},
{"ip_cidr": [Constants.SINGBOX_INTERNAL_SUBNET], "action": "hijack-dns"},
{"ip_cidr": [f"{server_ip}/32"], "outbound": "direct"},
{"ip_is_private": True, "outbound": "direct"},
{"ip_version": 6, "outbound": "block"},
{"inbound": ["tun-in", "socks-in"], "outbound": "proxy"},
],
"final": "proxy",
"default_domain_resolver": "tunnel-dns",
"auto_detect_interface": True,
},
}
def _cleanup_all() -> None:
try:
SingboxRunner().stop()
except Exception:
pass
try:
revert_dns_on_tun()
except Exception:
pass
try:
killswitch.disarm()
except Exception:
pass
time.sleep(1)
def enable_hysteria(username: str, password: str, server_host: str,
server_ip: str, socks5_port: int, observer=None) -> bool:
_cleanup_all()
runner = SingboxRunner()
config_path = Path(Constants.SINGBOX_CONFIG_DIR) / f"{username}-sing-box.json"
config = build_hysteria_config(username, password, server_host,
socks5_port, server_ip)
_connected = False
try:
runner.write_config(config_path, config)
if not runner.start(config_path):
if observer:
observer.notify("error", "sing-box not active after start")
return False
if not wait_for_tun(timeout=15.0):
if observer:
observer.notify("error", f"{Constants.SINGBOX_TUN_IF} did not appear after 15s")
return False
if not killswitch.arm(server_ip, Constants.SINGBOX_TUN_IF, Constants.SINGBOX_INTERNAL_SUBNET):
if observer:
observer.notify("error", "Failed to arm kill switch")
return False
# Fase 4 — DNS
_C = Constants()
dns_ok = set_dns_on_tun() if _C.HYSTERIA2_DNS_ENABLED else False
_connected = True
if observer:
observer.notify("connected", {
"tunnel_if": Constants.SINGBOX_TUN_IF,
"socks5_port": socks5_port,
"server_ip": server_ip,
"dns_enabled": _C.HYSTERIA2_DNS_ENABLED,
"dns_active": dns_ok,
})
return True
except Exception as e:
if observer:
observer.notify("error", str(e))
return False
finally:
if not _connected:
_cleanup_all()
def _force_kill_singbox() -> None:
import subprocess
try:
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
if result.returncode == 0:
for pid in result.stdout.strip().splitlines():
subprocess.run(
['sudo', 'kill', '-9', pid.strip()],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=5
)
except Exception:
pass
def disable_hysteria(observer=None) -> bool:
try:
revert_dns_on_tun()
except Exception:
pass
try:
SingboxRunner().stop()
except Exception:
pass
try:
import subprocess, time
time.sleep(0.5)
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
if result.returncode == 0:
_force_kill_singbox()
except Exception:
pass
try:
killswitch.disarm()
except Exception:
pass
if observer:
observer.notify("disconnected", {"tunnel_if": Constants.SINGBOX_TUN_IF})
return True

View file

@ -0,0 +1,196 @@
from core.errors.logger import logger
from core.utils.encrypted_proxy import killswitch
from typing import Optional, Dict, Any
from dataclasses import dataclass
@dataclass
class ConnectionStateData:
tunnel_interface: str = "?"
server_ip: Optional[str] = None
socks5_port: str = "?"
dns_enabled: bool = True
dns_active: bool = False
session_token: Optional[str] = None
def create_ui_state() -> Dict[str, Any]:
return {
'spinner': None,
'monitor': None,
'timer_state': None,
'connection_data': ConnectionStateData(),
}
# ── Cleanup helpers ───────────────────────────────────────────────────────────
def cleanup_spinner(state: Dict[str, Any]) -> None:
if state['spinner']:
try:
state['spinner'].stop()
except Exception as e:
logger.error(f"Failed to stop spinner: {e}")
finally:
state['spinner'] = None
def cleanup_monitor(state: Dict[str, Any]) -> None:
if state['monitor'] is not None:
try:
state['monitor'].stop()
except Exception as e:
logger.error(f"Failed to stop monitor: {e}")
finally:
state['monitor'] = None
def cleanup_timer(state: Dict[str, Any]) -> None:
if state['timer_state'] and state['timer_state'].get('stop'):
try:
state['timer_state']['stop'].set()
except Exception as e:
logger.error(f"Failed to stop timer: {e}")
finally:
if state['timer_state']:
state['timer_state']['stop'] = None
def cleanup_reconnecting_spinner(state: Dict[str, Any]) -> None:
if state['timer_state'] and state['timer_state'].get('reconnecting'):
try:
state['timer_state']['reconnecting'].set()
except Exception as e:
logger.error(f"Failed to stop reconnecting spinner: {e}")
finally:
if state['timer_state']:
state['timer_state']['reconnecting'] = None
def cleanup_all(state: Dict[str, Any]) -> None:
cleanup_spinner(state)
cleanup_timer(state)
cleanup_reconnecting_spinner(state)
cleanup_monitor(state)
# ── Handler factory ───────────────────────────────────────────────────────────
def create_event_handlers(state: Dict[str, Any], drop_state: Dict, tunnel_state: Dict):
"""Factory: returns event handler closures bound to state."""
def _update_connection_data(event) -> ConnectionStateData:
"""Extract connection fields from event subject into a ConnectionStateData."""
d = event.subject or {}
return ConnectionStateData(
tunnel_interface=d.get('tunnel_if', "?"),
server_ip=d.get('server_ip'),
socks5_port=d.get('socks5_port', "?"),
dns_enabled=d.get('dns_enabled', True),
dns_active=d.get('dns_active', False),
)
def _update_ui_labels(conn: ConnectionStateData) -> None:
"""Refresh all CLI status labels from current connection data."""
from cli.ui import (
label_connected, label_killswitch, label_dns, label_ipv6, label_ipv4
)
label_connected(f"tunnel={conn.tunnel_interface} port={conn.socks5_port}")
label_killswitch(killswitch.status())
label_dns(enabled=conn.dns_enabled, active=conn.dns_active)
label_ipv6(blocked=True)
if conn.server_ip:
label_ipv4(server_ip=conn.server_ip, reachable=True)
def _update_monitor_port(event_data: Dict) -> None:
"""Update monitor with new SOCKS5 port."""
if state['monitor'] is not None:
state['monitor'].set_connection_data(socks5_port=event_data.get('socks5_port'))
def on_connecting(event) -> None:
"""Handle connecting event: show attempt spinner."""
try:
cleanup_timer(state)
if (state['timer_state'] and
state['timer_state'].get('reconnecting') is None and
state['timer_state'].get('_retrying')):
from cli.ui import connecting_spinner
state['timer_state']['reconnecting'] = connecting_spinner()
d = event.subject or {}
attempt = d.get("attempt_count", "?")
total = d.get("maximum_number_of_attempts", "?")
cleanup_spinner(state)
from cli.ui import Spinner
state['spinner'] = Spinner(f"Connecting... attempt {attempt}/{total}")
state['spinner'].start()
except Exception as e:
logger.error(f"Error in on_connecting: {e}", exc_info=True)
def on_connected(event) -> None:
"""Handle connected event: update labels, monitor port, and start session timer."""
try:
cleanup_spinner(state)
cleanup_reconnecting_spinner(state)
if state['timer_state']:
state['timer_state']['_retrying'] = False
state['connection_data'] = _update_connection_data(event)
_update_ui_labels(state['connection_data'])
event_data = event.subject or {}
_update_monitor_port(event_data)
if state['timer_state'] is not None:
from cli.ui import session_timer
state['timer_state']['stop'] = session_timer(
drop_state=drop_state,
tunnel_state=tunnel_state,
)
except Exception as e:
logger.error(f"Error in on_connected: {e}", exc_info=True)
def on_connected_token(event) -> None:
"""Handle session token event: pass token to monitor."""
try:
if state['monitor'] is not None:
token = (event.subject or {}).get('session_token')
if token:
state['monitor'].set_session_token(token)
state['connection_data'].session_token = token
except Exception as e:
logger.error(f"Error in on_connected_token: {e}", exc_info=True)
def on_disconnected(event) -> None:
"""Handle disconnected event: cleanup UI and show disconnect label."""
try:
cleanup_all(state)
from cli.ui import label_disconnect
tunnel_if = (event.subject or {}).get('tunnel_if', '?')
label_disconnect(f"tunnel={tunnel_if}")
except Exception as e:
logger.error(f"Error in on_disconnected: {e}", exc_info=True)
def on_error(event) -> None:
"""Handle error event: cleanup UI and show error label."""
try:
cleanup_all(state)
from cli.ui import label_error
msg = event.subject
if isinstance(msg, dict):
msg = msg.get('message', str(msg))
label_error(str(msg or "Unknown error"))
logger.error(f"Connection error: {msg}")
except Exception as e:
logger.error(f"Error in on_error: {e}", exc_info=True)
return {
'on_connecting': on_connecting,
'on_connected': on_connected,
'on_connected_token': on_connected_token,
'on_disconnected': on_disconnected,
'on_error': on_error,
}

View file

@ -0,0 +1,203 @@
from urllib.parse import unquote
from pathlib import Path
import socket
import time
from core.Constants import Constants
from core.utils.encrypted_proxy.singbox import SingboxRunner
from core.utils.encrypted_proxy.dns import wait_for_tun, set_dns_on_tun, revert_dns_on_tun
from core.utils.encrypted_proxy import killswitch
def parse_vless_link(link: str) -> dict:
link = link.replace("vless://", "")
uuid, rest = link.split("@", 1)
hostport, qs = rest.split("?", 1)
query = qs.split("#")[0]
host, port = hostport.rsplit(":", 1)
params = {}
for part in query.split("&"):
if "=" in part:
k, v = part.split("=", 1)
params[k] = v
sni = params.get("sni", host)
ws_host = params.get("host", "").strip() or sni
return {
"uuid": uuid,
"host": host,
"port": int(port),
"path": unquote(params.get("path", "/vless")),
"sni": sni,
"ws_host": ws_host,
"security": params.get("security", "tls"),
"network": params.get("type", "ws"),
}
def build_vless_config(vless: dict, socks5_port: int, server_ip: str) -> dict:
return {
"dns": {
"servers": [{"tag": "tunnel-dns", "type": "udp", "server": "9.9.9.9"}],
"final": "tunnel-dns",
"strategy": "ipv4_only",
"independent_cache": True,
},
"inbounds": [
{
"type": "tun",
"tag": "tun-in",
"interface_name": Constants.SINGBOX_TUN_IF,
"address": [Constants.SINGBOX_INTERNAL_ADDR],
"mtu": 9000,
"auto_route": True,
"stack": "gvisor",
},
{
"type": "socks",
"tag": "socks-in",
"listen": "127.0.0.1",
"listen_port": socks5_port,
},
],
"outbounds": [
{"type": "direct", "tag": "direct"},
{"type": "block", "tag": "block"},
{
"type": "vless",
"tag": "proxy",
"server": server_ip,
"server_port": vless["port"],
"uuid": vless["uuid"],
"tls": {
"enabled": vless["security"] == "tls",
"server_name": vless["sni"],
"insecure": False,
},
"transport": {
"type": "ws",
"path": vless["path"],
"headers": {"Host": vless["ws_host"]},
},
},
],
"route": {
"rules": [
{"protocol": "dns", "action": "hijack-dns"},
{"ip_cidr": [Constants.SINGBOX_INTERNAL_SUBNET], "action": "hijack-dns"},
{"ip_cidr": [f"{server_ip}/32"], "outbound": "direct"},
{"ip_is_private": True, "outbound": "direct"},
{"ip_version": 6, "outbound": "block"},
{"inbound": ["tun-in", "socks-in"], "outbound": "proxy"},
],
"final": "proxy",
"default_domain_resolver": "tunnel-dns",
"auto_detect_interface": True,
},
}
def _cleanup_all() -> None:
try:
SingboxRunner().stop()
except Exception:
pass
try:
revert_dns_on_tun()
except Exception:
pass
try:
killswitch.disarm()
except Exception:
pass
time.sleep(1)
def enable_vless(vless_link: str, username: str, server_ip: str,
socks5_port: int, observer=None) -> bool:
_cleanup_all()
vless = parse_vless_link(vless_link)
runner = SingboxRunner()
config_path = Path(Constants.SINGBOX_CONFIG_DIR) / f"{username}-sing-box.json"
config = build_vless_config(vless, socks5_port, server_ip)
_connected = False
try:
runner.write_config(config_path, config)
if not runner.start(config_path):
if observer:
observer.notify("error", "sing-box not active after start")
return False
if not wait_for_tun(timeout=15.0):
if observer:
observer.notify("error", f"{Constants.SINGBOX_TUN_IF} did not appear after 15s")
return False
if not killswitch.arm(server_ip, Constants.SINGBOX_TUN_IF, Constants.SINGBOX_INTERNAL_SUBNET):
if observer:
observer.notify("error", "Failed to arm kill switch")
return False
_C = Constants()
dns_ok = set_dns_on_tun() if _C.VLESS_DNS_ENABLED else False
_connected = True
if observer:
observer.notify("connected", {
"tunnel_if": Constants.SINGBOX_TUN_IF,
"socks5_port": socks5_port,
"server_ip": server_ip,
"dns_enabled": _C.VLESS_DNS_ENABLED,
"dns_active": dns_ok,
})
return True
except Exception as e:
if observer:
observer.notify("error", str(e))
return False
finally:
if not _connected:
_cleanup_all()
def _force_kill_singbox() -> None:
import subprocess
try:
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
if result.returncode == 0:
for pid in result.stdout.strip().splitlines():
subprocess.run(
['sudo', 'kill', '-9', pid.strip()],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, timeout=5
)
except Exception:
pass
def disable_vless(observer=None) -> bool:
try:
revert_dns_on_tun()
except Exception:
pass
try:
SingboxRunner().stop()
except Exception:
pass
try:
import subprocess, time
time.sleep(0.5)
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
if result.returncode == 0:
_force_kill_singbox()
except Exception:
pass
try:
killswitch.disarm()
except Exception:
pass
if observer:
observer.notify("disconnected", {"tunnel_if": Constants.SINGBOX_TUN_IF})
return True

View file

View file

@ -0,0 +1,72 @@
import os
import subprocess
import time
from core.Constants import Constants
_C = Constants()
_SUDO_KW = dict(
stdout=subprocess.DEVNULL,
stderr=subprocess.PIPE,
stdin=subprocess.DEVNULL,
env={**os.environ, "SUDO_ASKPASS": "/bin/false"},
text=True,
)
def is_resolved_active() -> bool:
try:
result = subprocess.run(
["systemctl", "is-active", "systemd-resolved"],
capture_output=True,
text=True,
timeout=5,
)
return result.stdout.strip() == "active"
except (subprocess.TimeoutExpired, FileNotFoundError):
return False
def wait_for_tun(timeout: float = 15.0, interval: float = 0.5) -> bool:
elapsed = 0.0
while elapsed < timeout:
result = subprocess.run(
["ip", "link", "show", _C.SINGBOX_TUN_IF],
capture_output=True,
)
if result.returncode == 0:
return True
time.sleep(interval)
elapsed += interval
return False
def set_dns_on_tun() -> bool:
if not is_resolved_active():
print("[DNS] systemd-resolved is not active — DNS on tun skipped")
return False
result = subprocess.run(
["sudo", _C.RESOLVECTL_WRAPPER, "set", "9.9.9.9"],
**_SUDO_KW,
)
if result.returncode != 0:
print(f"[DNS] Failed to configure {_C.SINGBOX_TUN_IF}: {result.stderr.strip()}")
return False
return True
def revert_dns_on_tun() -> bool:
if not is_resolved_active():
return True
result = subprocess.run(
["sudo", _C.RESOLVECTL_WRAPPER, "revert"],
**_SUDO_KW,
)
if result.returncode != 0:
print(f"[DNS] Failed to revert {_C.SINGBOX_TUN_IF}: {result.stderr.strip()}")
return False
return True

View file

@ -0,0 +1,36 @@
import subprocess
from core.Constants import Constants
_C = Constants()
def arm(server_ip: str, tunnel_if: str, internal_subnet: str = None) -> bool:
cmd = ["sudo", _C.KILLSWITCH_WRAPPER, "arm", server_ip, tunnel_if]
if internal_subnet:
cmd.append(internal_subnet)
result = subprocess.run(cmd, capture_output=True, text=True)
if result.returncode != 0:
print(f"[killswitch] Failed to arm: {result.stderr.strip()}")
return False
return True
def disarm() -> bool:
result = subprocess.run(
["sudo", _C.KILLSWITCH_WRAPPER, "disarm"],
capture_output=True,
text=True,
)
if result.returncode != 0:
print(f"[killswitch] Failed to disarm: {result.stderr.strip()}")
return False
return True
def status() -> bool:
result = subprocess.run(
["sudo", _C.KILLSWITCH_WRAPPER, "status"],
capture_output=True,
text=True,
)
return result.returncode == 0 and result.stdout.strip() == "armed"

View file

@ -0,0 +1,147 @@
import subprocess
import threading
import time
from core.utils.encrypted_proxy import killswitch
from core.utils.encrypted_proxy.singbox import SingboxRunner, SingboxProcessMonitor
from core.utils.encrypted_proxy.singbox import CAUSE_DISCONNECT, CAUSE_LEAK, CAUSE_DISPLACED
drop_state: dict = {"count": 0}
tunnel_state: dict = {"alive": True}
class KillswitchMonitor(SingboxProcessMonitor):
"""CLI-specific monitor for sing-box (vless/hysteria2): checks sing-box process health."""
LOG_PREFIX = "hydraveil-drop"
def __init__(self):
super().__init__()
self._runner = SingboxRunner()
self._on_disconnect = None
self._on_leak = None
self._on_displaced = None
self._socks5_port = None
self._session_token = None
self._cause_event = threading.Event()
self._cause: str | None = None
# ── Setters ───────────────────────────────────────────────────────────────
def set_on_disconnect(self, fn):
self._on_disconnect = fn
def set_on_leak(self, fn):
self._on_leak = fn
def set_on_displaced(self, fn):
self._on_displaced = fn
def set_connection_data(self, socks5_port: int):
self._socks5_port = socks5_port
def set_session_token(self, token: str):
self._session_token = token
# ── Abstract impl ─────────────────────────────────────────────────────────
def is_tunnel_alive(self) -> bool:
return self._runner.is_running()
def on_process_alive(self) -> None:
tunnel_state["alive"] = True
def on_process_dead(self) -> None:
tunnel_state["alive"] = False
def on_monitor_stopped(self, cause: str) -> None:
if cause == CAUSE_DISCONNECT and self._on_disconnect:
self._on_disconnect()
elif cause == CAUSE_LEAK and self._on_leak:
self._on_leak()
elif cause == CAUSE_DISPLACED and self._on_displaced:
self._on_displaced()
# ── Override monitor loop to add CLI-specific checks ──────────────────────
def _monitor_loop(self) -> None:
while True:
with self._lock:
if not self._running:
break
interrupted = self._stop_event.wait(timeout=self.POLL_INTERVAL)
if interrupted:
break
# Killswitch check
if not killswitch.status():
tunnel_state["alive"] = False
self._trigger(CAUSE_DISCONNECT)
return
# Sing-box check
alive = self._runner.is_running()
tunnel_state["alive"] = alive
if not alive:
self._trigger(CAUSE_LEAK)
return
# Session token check
if self._session_token is not None:
from core.models.system.SystemState import SystemState
current = SystemState.get()
if current is None or current.session_token != self._session_token:
tunnel_state["alive"] = False
self._trigger(CAUSE_DISPLACED)
return
self._poll_drops()
def _trigger(self, cause: str):
self._cause = cause
with self._lock:
self._running = False
self._cause_event.set()
self._stop_event.set()
def _poll_drops(self) -> None:
since = self._last_poll_ts or "10 seconds ago"
try:
result = subprocess.run(
[
"journalctl", "-k", "--no-pager",
f"--grep={self.LOG_PREFIX}",
"--since", since,
"-o", "short-iso",
],
capture_output=True,
text=True,
timeout=5,
)
new_drops = result.stdout.count(self.LOG_PREFIX)
if new_drops:
drop_state["count"] += new_drops
except (subprocess.TimeoutExpired, FileNotFoundError):
pass
finally:
self._last_poll_ts = time.strftime("%Y-%m-%dT%H:%M:%S")
# ── Override run_blocking to handle cause_event ───────────────────────────
def run_blocking(self) -> None:
with self._lock:
self._running = True
self._stop_event.clear()
self._cause_event.clear()
self._cause = None
drop_state["count"] = 0
tunnel_state["alive"] = True
self._last_poll_ts = time.strftime("%Y-%m-%dT%H:%M:%S")
t = threading.Thread(target=self._monitor_loop, daemon=True)
t.start()
self._cause_event.wait()
t.join(timeout=2)
self.on_monitor_stopped(self._cause or "user_requested")

View file

@ -0,0 +1,74 @@
import socket
import subprocess
import time
def resolve(host: str) -> str | None:
try:
return socket.gethostbyname(host)
except Exception:
return None
def get_real_ip(timeout: int = 10) -> str:
endpoints = [
"https://api.ipify.org",
"https://ifconfig.me/ip",
"https://icanhazip.com",
]
for url in endpoints:
try:
r = subprocess.run(
["curl", "-s", "--max-time", str(timeout), url],
capture_output=True,
timeout=timeout + 2,
)
ip = r.stdout.decode().strip()
if ip and not ip.startswith("unknown"):
return ip
except Exception:
continue
return "unknown"
def get_proxied_ip(socks5_port: int, timeout: int = 15) -> str:
endpoints = [
"https://api.ipify.org",
"https://ifconfig.me/ip",
"https://icanhazip.com",
]
for url in endpoints:
try:
r = subprocess.run(
[
"curl", "-s",
"--max-time", str(timeout),
"--connect-timeout", "10",
"-x", f"socks5h://127.0.0.1:{socks5_port}",
url,
],
capture_output=True,
timeout=timeout + 2,
)
ip = r.stdout.decode().strip()
stderr = r.stderr.decode().strip()
if ip and not ip.startswith("<") and "." in ip:
return ip
if stderr:
print(f"[net] curl stderr ({url}): {stderr[:120]}")
except subprocess.TimeoutExpired:
print(f"[net] timeout via socks5 → {url}")
except Exception as e:
print(f"[net] error via socks5 → {url}: {e}")
return "unknown"
def verify_ip(socks5_port: int, retries: int = 5, delay: float = 3.0) -> str:
for attempt in range(1, retries + 1):
print(f"[net] verify_ip attempt {attempt}/{retries}...")
ip = get_proxied_ip(socks5_port)
if ip != "unknown":
return ip
if attempt < retries:
time.sleep(delay)
return "unknown"

View file

@ -0,0 +1,335 @@
import json
import os
import subprocess
import time
from abc import ABC, abstractmethod
from pathlib import Path
from typing import Optional
import threading
from core.Constants import Constants
import logging
logger = logging.getLogger(__name__)
class SingboxRunner:
_WRAPPER = Constants.SINGBOX_WRAPPER
_PID_FILE = Path(Constants.SINGBOX_PID_FILE)
_LOG_FILE = Path(Constants.SINGBOX_LOG_FILE)
_CONFIG_DIR = Path(Constants.SINGBOX_CONFIG_DIR)
_START_TIMEOUT_S = 15.0
_START_POLL_S = 0.5
_START_INITIAL_S = 2.0
_ORPHAN_GRACEFUL_TIMEOUT_S = 2.0
_ORPHAN_FINAL_VERIFY_TIMEOUT_S = 1.0
_ORPHAN_POLL_S = 0.3
def __init__(self):
self._wrapper_process: Optional[subprocess.Popen] = None
def start(self, config_path: Path) -> bool:
self.stop()
self._assert_config_safe(config_path)
self._ensure_log_file()
try:
self._wrapper_process = subprocess.Popen(
['sudo', self._WRAPPER, str(config_path), 'start'],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stdin=subprocess.DEVNULL,
env={**os.environ, 'SUDO_ASKPASS': '/bin/false'},
)
logger.info(f"Wrapper spawned (PID {self._wrapper_process.pid})")
except FileNotFoundError:
raise RuntimeError(
f'Wrapper not found: {self._WRAPPER}\n'
'Run the installer first.'
)
time.sleep(self._START_INITIAL_S)
deadline = time.monotonic() + self._START_TIMEOUT_S
while time.monotonic() < deadline:
if self.is_running():
return True
time.sleep(self._START_POLL_S)
raise RuntimeError(
f'sing-box did not start within {self._START_TIMEOUT_S}s.\n'
f'Last error:\n{self.get_last_error()}'
)
def stop(self) -> None:
if self._wrapper_process is not None:
try:
self._wrapper_process.terminate()
self._wrapper_process.wait(timeout=2)
except subprocess.TimeoutExpired:
self._wrapper_process.kill()
finally:
self._wrapper_process = None
self._kill_orphans()
try:
subprocess.run(
['sudo', self._WRAPPER, '', 'stop'],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stdin=subprocess.DEVNULL,
env={**os.environ, 'SUDO_ASKPASS': '/bin/false'},
timeout=5,
)
except (subprocess.TimeoutExpired, FileNotFoundError, OSError):
logger.warning("Graceful stop failed, attempting emergency kill")
self._emergency_kill()
def is_running(self) -> bool:
if not self._PID_FILE.exists():
try:
result = subprocess.run(
['pgrep', '-x', 'sing-box'],
capture_output=True,
text=True,
)
return result.returncode == 0
except OSError:
return False
try:
pid = int(self._PID_FILE.read_text().strip())
if pid <= 0:
raise ValueError
os.kill(pid, 0)
return True
except (ValueError, ProcessLookupError):
self._PID_FILE.unlink(missing_ok=True)
return False
except PermissionError:
return True
def get_last_error(self) -> str:
try:
result = subprocess.run(
['sudo', self._WRAPPER, '', 'log'],
capture_output=True,
text=True,
timeout=5,
stdin=subprocess.DEVNULL,
env={**os.environ, 'SUDO_ASKPASS': '/bin/false'},
)
return result.stdout.strip()
except (subprocess.TimeoutExpired, FileNotFoundError, OSError):
try:
lines = self._LOG_FILE.read_text(errors='replace').splitlines()
return '\n'.join(lines[-80:])
except OSError:
return ''
def write_config(self, config_path: Path, config: dict) -> None:
self._assert_config_safe(config_path)
config_path.parent.mkdir(parents=True, exist_ok=True, mode=0o700)
tmp = config_path.with_suffix('.tmp')
try:
with open(tmp, 'w') as f:
json.dump(config, f, indent=2)
os.chmod(tmp, 0o600)
tmp.rename(config_path)
except Exception:
tmp.unlink(missing_ok=True)
raise
# ── Kill helpers ──────────────────────────────────────────────────────────
def _kill_orphans(self) -> None:
"""Kill any lingering sing-box processes using escalating force."""
# Phase 1: Check if anything is actually running
if not self._has_running_processes():
return
logger.debug("Found running sing-box processes, attempting graceful stop")
# Phase 2: Graceful stop attempt via wrapper
self._attempt_graceful_stop()
# Phase 3: Wait for graceful shutdown to complete
if self._wait_for_process_exit(self._ORPHAN_GRACEFUL_TIMEOUT_S):
logger.info("sing-box stopped gracefully")
return
# Phase 4: Graceful stop failed — force kill remaining processes
logger.warning("Graceful stop timeout, force killing remaining sing-box processes")
self._force_kill_all_processes()
# Phase 5: Verify all processes are actually dead
if self._wait_for_process_exit(self._ORPHAN_FINAL_VERIFY_TIMEOUT_S):
logger.info("All sing-box processes terminated")
return
logger.error("Failed to terminate all sing-box processes after force kill")
def _has_running_processes(self) -> bool:
try:
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
return result.returncode == 0
except OSError:
return False
def _attempt_graceful_stop(self) -> None:
try:
subprocess.run(
['sudo', self._WRAPPER, '', 'stop'],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stdin=subprocess.DEVNULL,
env={**os.environ, 'SUDO_ASKPASS': '/bin/false'},
timeout=5,
)
except (OSError, subprocess.TimeoutExpired):
pass
def _wait_for_process_exit(self, timeout_s: float) -> bool:
deadline = time.monotonic() + timeout_s
while time.monotonic() < deadline:
if not self._has_running_processes():
return True
time.sleep(self._ORPHAN_POLL_S)
return False
def _force_kill_all_processes(self) -> None:
try:
result = subprocess.run(['pgrep', '-x', 'sing-box'], capture_output=True, text=True)
if result.returncode != 0:
return
pids = [int(p) for p in result.stdout.strip().splitlines() if p.strip().isdigit()]
for pid in pids:
self._force_kill_process(pid)
except (OSError, subprocess.TimeoutExpired):
pass
def _force_kill_process(self, pid: int) -> None:
try:
subprocess.run(
['sudo', 'kill', '-9', str(pid)],
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
stdin=subprocess.DEVNULL,
env={**os.environ, 'SUDO_ASKPASS': '/bin/false'},
timeout=5,
)
logger.debug(f"Force killed sing-box process {pid}")
except (OSError, subprocess.TimeoutExpired):
pass
def _emergency_kill(self) -> None:
if not self._PID_FILE.exists():
return
try:
pid = int(self._PID_FILE.read_text().strip())
if pid > 0:
self._force_kill_process(pid)
except (ValueError, OSError):
pass
finally:
self._PID_FILE.unlink(missing_ok=True)
def _ensure_log_file(self) -> None:
self._LOG_FILE.parent.mkdir(parents=True, exist_ok=True)
if not self._LOG_FILE.exists():
self._LOG_FILE.touch(mode=0o600)
def _assert_config_safe(self, config_path: Path) -> None:
try:
config_path.resolve().relative_to(self._CONFIG_DIR.resolve())
except ValueError:
raise ValueError(
f'Config outside the allowed directory.\n'
f' Config: {config_path}\n'
f' Allowed: {self._CONFIG_DIR}'
)
# ── Abstract base monitor ─────────────────────────────────────────────────────
CAUSE_DISCONNECT = 'disconnect'
CAUSE_LEAK = 'leak'
CAUSE_DISPLACED = 'displaced'
class SingboxProcessMonitor(ABC):
"""Abstract monitor for sing-box process health.
Clients inherit this and implement event callbacks based on their
threading model (threading.Event for CLI, Qt signals for GUI).
"""
POLL_INTERVAL = 10
def __init__(self):
self._lock = threading.Lock()
self._running = False
self._stop_event = threading.Event()
self._last_poll_ts = None
def stop(self):
with self._lock:
self._running = False
self._stop_event.set()
@abstractmethod
def is_tunnel_alive(self) -> bool:
"""Return True if the tunnel is healthy."""
pass
@abstractmethod
def on_process_alive(self) -> None:
pass
@abstractmethod
def on_process_dead(self) -> None:
pass
@abstractmethod
def on_monitor_stopped(self, cause: str) -> None:
pass
def _monitor_loop(self) -> None:
while True:
with self._lock:
if not self._running:
break
interrupted = self._stop_event.wait(timeout=self.POLL_INTERVAL)
if interrupted:
break
alive = self.is_tunnel_alive()
if alive:
self.on_process_alive()
else:
self.on_process_dead()
break
def run_blocking(self) -> None:
with self._lock:
self._running = True
self._stop_event.clear()
t = threading.Thread(target=self._monitor_loop, daemon=False)
t.start()
self._stop_event.wait()
with self._lock:
self._running = False
t.join(timeout=5)
# tunnel teardown handled by subclass
self.on_monitor_stopped("user_requested")

View file

View file

@ -0,0 +1,104 @@
import subprocess
import threading
from core.utils.encrypted_proxy.singbox import SingboxProcessMonitor
from core.utils.encrypted_proxy.singbox import CAUSE_DISCONNECT, CAUSE_DISPLACED
class WireGuardMonitor(SingboxProcessMonitor):
"""Monitor for WireGuard system connections.
Checks if the wg interface is alive instead of a sing-box process.
Shared between CLI and GUI.
"""
def __init__(self):
super().__init__()
self._session_token = None
self._on_disconnect = None
self._on_displaced = None
self._cause_event = threading.Event()
self._cause: str | None = None
# ── Setters ───────────────────────────────────────────────────────────────
def set_on_disconnect(self, fn):
self._on_disconnect = fn
def set_on_displaced(self, fn):
self._on_displaced = fn
def set_session_token(self, token: str):
self._session_token = token
# ── Tunnel health check ───────────────────────────────────────────────────
def is_tunnel_alive(self) -> bool:
try:
result = subprocess.run(
['ip', 'link', 'show', 'wg'],
capture_output=True,
timeout=5,
)
return result.returncode == 0
except (subprocess.TimeoutExpired, FileNotFoundError):
return False
# ── Abstract impl ─────────────────────────────────────────────────────────
def on_process_alive(self) -> None:
pass
def on_process_dead(self) -> None:
self._trigger(CAUSE_DISCONNECT)
def on_monitor_stopped(self, cause: str) -> None:
if cause == CAUSE_DISCONNECT and self._on_disconnect:
self._on_disconnect()
elif cause == CAUSE_DISPLACED and self._on_displaced:
self._on_displaced()
# ── Monitor loop override ─────────────────────────────────────────────────
def _monitor_loop(self) -> None:
while True:
with self._lock:
if not self._running:
break
interrupted = self._stop_event.wait(timeout=self.POLL_INTERVAL)
if interrupted:
break
if not self.is_tunnel_alive():
self._trigger(CAUSE_DISCONNECT)
return
if self._session_token is not None:
from core.models.system.SystemState import SystemState
current = SystemState.get()
if current is None or current.session_token != self._session_token:
self._trigger(CAUSE_DISPLACED)
return
def _trigger(self, cause: str):
self._cause = cause
with self._lock:
self._running = False
self._cause_event.set()
self._stop_event.set()
# ── run_blocking override ─────────────────────────────────────────────────
def run_blocking(self) -> None:
with self._lock:
self._running = True
self._stop_event.clear()
self._cause_event.clear()
self._cause = None
t = threading.Thread(target=self._monitor_loop, daemon=True)
t.start()
self._cause_event.wait()
t.join(timeout=2)
self.on_monitor_stopped(self._cause or "user_requested")

View file

@ -13,6 +13,7 @@ classifiers = [
]
dependencies = [
"cryptography ~= 46.0.3",
"python-dotenv ~= 1.0.0",
"dataclasses-json ~= 0.6.7",
"marshmallow ~= 3.26.1",
"pysocks ~= 1.7.1",