diff --git a/change_log.md b/change_log.md
index c393f42..364a469 100644
--- a/change_log.md
+++ b/change_log.md
@@ -1,5 +1,10 @@
# Major Change Log:
+# EncryptedProxy Configs
+### Aug 6, 2026
+EncryptedProxy configuration path setup, but not yet functional. Pydantic models for each protocol are being setup.
+
+
# Transition
### Aug 6, 2026
1) Transition ticket prep from old requests-based GET/POST system to the new httpx one.
diff --git a/core/models/Result.py b/core/models/Result.py
index 3178e27..b8766ba 100644
--- a/core/models/Result.py
+++ b/core/models/Result.py
@@ -12,6 +12,7 @@ class ResultError(Enum):
MISSING_FILE = "missing_file"
MISSING_DEPENDENCY = "missing_dependency"
MISSING_DATA = "missing_data"
+ FILE_SYSTEM = "filesystem"
CONNECTION = "connection"
DATABASE = "database"
PERMISSION = "permission"
@@ -24,6 +25,7 @@ class ResultError(Enum):
EXTERNAL_DNS = "external_dns"
INTERFACE = "interface"
TIMEOUT = "timeout"
+ INVALID_API_REPLY = "invalid_api_reply"
UNKNOWN = "unknown"
@dataclass
diff --git a/core/models/manage/pydantic_manager.py b/core/models/manage/pydantic_manager.py
new file mode 100644
index 0000000..1f65dbb
--- /dev/null
+++ b/core/models/manage/pydantic_manager.py
@@ -0,0 +1,96 @@
+from core.models.HysteriaConfig import HysteriaConfig
+from core.errors.logger import logger
+
+# generic
+import os
+import json
+from pathlib import Path
+from pydantic import BaseModel
+import subprocess
+
+
+def save_to_sudo_folder(model: BaseModel, filepath: str) -> bool:
+ """
+ Save model to a sudo-protected file using pkexec.
+ """
+ try:
+ Path(filepath).parent.mkdir(parents=True, exist_ok=True)
+
+ json_data = model.model_dump_json(indent=2)
+
+ # Use pkexec + tee to write with elevated privileges
+ process = subprocess.Popen(
+ ('pkexec', 'tee', filepath),
+ stdin=subprocess.PIPE,
+ stdout=subprocess.PIPE,
+ stderr=subprocess.PIPE,
+ text=True
+ )
+ stdout, stderr = process.communicate(input=json_data)
+
+ if process.returncode != 0:
+ logger.error(f"Error: {stderr}")
+ return False
+
+ logger.info(f"Saved to {filepath} (elevated)")
+ return True
+
+ except FileNotFoundError as e:
+ logger.error(f"Error: Parent directory doesn't exist: {filepath}")
+ return False
+
+ except PermissionError as e:
+ logger.error(f"Error: Permission denied writing to {filepath}")
+ return False
+
+ except OSError as e:
+ logger.error(f"Error: OS error (disk full?): {e}")
+ return False
+
+ except TypeError as e:
+ logger.error(f"Error: Model serialization failed (invalid type): {e}")
+ return False
+
+ except Exception as e:
+ logger.error(f"Error: Unexpected error saving {filepath}: {e}")
+ return False
+
+
+
+def save_to_regular_folder(model: BaseModel, filepath: str) -> bool:
+ """
+ Purpose:
+ Serialize a Pydantic model to JSON file WITHOUT sudo.
+
+ Returns:
+ Returns True on success,
+ False on failure.
+ """
+ try:
+ # Ensure parent directory exists
+ Path(filepath).parent.mkdir(parents=True, exist_ok=True)
+
+ with open(filepath, 'w') as f:
+ f.write(model.model_dump_json(indent=2))
+
+ return True
+
+ except FileNotFoundError as e:
+ logger.error(f"Error: Parent directory doesn't exist: {filepath}")
+ return False
+
+ except PermissionError as e:
+ logger.error(f"Error: Permission denied writing to {filepath}")
+ return False
+
+ except OSError as e:
+ logger.error(f"Error: OS error (disk full?): {e}")
+ return False
+
+ except TypeError as e:
+ logger.error(f"Error: Model serialization failed (invalid type): {e}")
+ return False
+
+ except Exception as e:
+ logger.error(f"Error: Unexpected error saving {filepath}: {e}")
+ return False
diff --git a/core/models/pydantic_models/HysteriaData.py b/core/models/pydantic_models/HysteriaData.py
new file mode 100644
index 0000000..ed46066
--- /dev/null
+++ b/core/models/pydantic_models/HysteriaData.py
@@ -0,0 +1,42 @@
+# generic
+from pydantic import BaseModel, field_validator, ValidationError, HttpUrl
+from sqlalchemy.orm import Session
+from pydantic import model_validator, ValidationInfo
+from typing_extensions import Self
+from ipaddress import IPv4Address
+from pydantic_core import PydanticUndefinedType
+
+class HysteriaData(BaseModel):
+ model_config = ConfigDict(extra="allow")
+
+ username: str
+ password: str
+ operator_hysteria2_host: HttpUrl
+ # operator_id: int
+ server_ip: IPv4Address
+ location_country_code: str
+ location_city_code: str
+
+ @model_validator(mode='before')
+ @classmethod
+ def denormalize(cls, data):
+ if isinstance(data, dict) and 'operator' in data:
+ return {
+ 'server_ip': data['operator'].get('id'),
+ 'api_url': data['operator'].get('domain'),
+ 'operator_hysteria2_host': data['operator'].get('hysteria2_host'),
+ }
+ return data
+
+ # @field_validator('operator_id')
+ # @classmethod
+ # def validate_operator_exists(cls, v, info):
+ # db = info.context.get('db')
+ # if not db:
+ # raise ValueError("Database session not provided")
+
+ # operator = db.query(Operator).filter(Operator.id == v).first()
+ # if not operator:
+ # raise ValueError(f"Operator ID ID {v} does not exist")
+
+ # return v
\ No newline at end of file
diff --git a/core/services/networking/general_connection_tools/connection_enable.py b/core/services/networking/general_connection_tools/connection_enable.py
index 51e094f..2d873cd 100644
--- a/core/services/networking/general_connection_tools/connection_enable.py
+++ b/core/services/networking/general_connection_tools/connection_enable.py
@@ -2,6 +2,7 @@ from core.services.networking.general_connection_tools.testing_evaluating import
from core.services.networking.systemwide.systemwide_wireguard import establish_system_connection, terminate_system_connection
from core.services.keys_and_verifications.wireguard_keys import register_wireguard_session
from core.services.subscriptions.subscriptions import activate_subscription
+# from core.services.networking.systemwide.encrypted_proxy.configure_singbox import configure_singbox
# If refactored to enums:
# from core.models.session.SessionConnection import SessionConnectionTypes
@@ -22,6 +23,8 @@ from core.models.BaseProfile import ProfileType
from core.observers.ConnectionObserver import ConnectionObserver
from core.controllers.SystemStateController import SystemStateController
+from core.services.networking.systemwide.encrypted_proxy.ensure_singbox_configured import ensure_singbox_configured
+
def establish_connection(
profile: Union[SessionProfile, SystemProfile],
ignore: tuple[type[Exception]] = (),
@@ -32,12 +35,23 @@ def establish_connection(
logger.info(f"[CONNECTION] Checking subscription..")
activate_subscription(profile, connection_observer)
+ # =========================================
+ # HYSTERIA2 & VLESS
+ # =========================================
+ # if profile.connection.code in ("hysteria2", "vless"):
+ # logger.info("Pulling encrypted proxies off the main flow..")
+ # configure_singbox(profile, connection_observer)
+
+ # =========================================
+ # SOCKS5 & WIREGUARD
+ # =========================================
logger.info(f"[CONNECTION] Checking proxy configuration..")
_ensure_proxy_configured(profile, connection_observer)
logger.info(f"[CONNECTION] Checking Wireguard configuration..")
- _ensure_wireguard_configured(profile, connection_observer)
-
+ if profile.connection.code == "wireguard":
+ _ensure_wireguard_configured(profile, connection_observer)
+
establish_fn = {
ProfileType.SESSION: ConnectionController.establish_session_connection,
ProfileType.SYSTEM: establish_system_connection,
@@ -131,6 +145,7 @@ def __should_renegotiate(profile: Union[SessionProfile, SystemProfile]):
+
# Enums
# if profile.type == ProfileType.SYSTEM:
# print("This is a system profile")
diff --git a/core/services/networking/httpx/connect.py b/core/services/networking/httpx/connect.py
index a1d4c72..b45c241 100644
--- a/core/services/networking/httpx/connect.py
+++ b/core/services/networking/httpx/connect.py
@@ -26,7 +26,7 @@ import httpx
_port_used = None
-def single_endpoint(method: str, url: str, observer: ConnectionObserver, payload: dict = None) -> ApiResponse:
+def single_endpoint(method: str, url: str, observer: ConnectionObserver, payload: dict = None, billing_code: str = None) -> ApiResponse:
"""
Rank:
Orchestrator
@@ -66,7 +66,8 @@ def single_endpoint(method: str, url: str, observer: ConnectionObserver, payload
method=method,
url=url,
client=client,
- payload=payload
+ payload=payload,
+ billing_code=billing_code
)
if initial_result.valid:
return initial_result
@@ -182,7 +183,8 @@ def bootstrap_and_try_again(method: str, url: str, observer: ConnectionObserver,
method=method,
url=url,
client=client,
- payload=payload
+ payload=payload,
+ billing_code=billing_code
)
@@ -335,7 +337,7 @@ def custom_dns_resolver_for_BULK_THREADING(dns_problem_list: list, observer: Con
return parallel_thread(client, desired_endpoints=dns_endpoints) # threading returns a dict!
-def custom_dns_resolver_for_SINGLE_ENDPOINT(method: str, url: str, observer: ConnectionObserver, payload: dict = None) -> dict:
+def custom_dns_resolver_for_SINGLE_ENDPOINT(method: str, url: str, observer: ConnectionObserver, payload: dict = None, billing_code: str = None) -> dict:
########################################################
# MAKE DNS RESOLVER
########################################################
@@ -362,7 +364,8 @@ def custom_dns_resolver_for_SINGLE_ENDPOINT(method: str, url: str, observer: Con
method=method,
url=url_with_ip,
client=client,
- payload=payload
+ payload=payload,
+ billing_code=billing_code
)
diff --git a/core/services/networking/httpx/make_request.py b/core/services/networking/httpx/make_request.py
index 6695c6d..3ba97fc 100644
--- a/core/services/networking/httpx/make_request.py
+++ b/core/services/networking/httpx/make_request.py
@@ -15,6 +15,7 @@ def make_request(
url: str,
client: httpx.Client,
payload: Optional[dict] = None,
+ billing_code: Optional[str] = None,
) -> ApiResponse:
if method == "post" and not payload:
@@ -39,8 +40,10 @@ def make_request(
method=method,
url=url,
client=client,
- payload=payload
- )
+ payload=payload,
+ billing_code=billing_code
+ ),
+
if second_result.valid:
return second_result
@@ -72,13 +75,27 @@ def _make_request(
url: str,
client: httpx.Client,
payload: Optional[dict] = None,
+ billing_code: Optional[str] = None,
) -> ApiResponse:
logger.debug(f"Executing {method.upper()} to {url}")
+
+ # ========== HEADERS ==========
+ if billing_code:
+ custom_headers = {'X-Billing-Code': billing_code}
+ else:
+ custom_headers = None
+
try:
+ # ========== GET ==========
if method.lower() == "get":
- response = client.get(url)
+ response = client.get(url, headers=custom_headers)
+ # ========== POST ==========
else:
- response = client.post(url, json=payload)
+ response = client.post(
+ url,
+ json=payload,
+ headers=custom_headers
+ )
return classify_response(response)
diff --git a/core/services/networking/systemwide/encrypted_proxy/configure_singbox.py b/core/services/networking/systemwide/encrypted_proxy/configure_singbox.py
new file mode 100644
index 0000000..ba6ef18
--- /dev/null
+++ b/core/services/networking/systemwide/encrypted_proxy/configure_singbox.py
@@ -0,0 +1,142 @@
+from core.models.orm_calls.location_calls import get_profile_location_data
+from core.services.networking.tor_tools import ports
+from core.services.networking.httpx import connect
+from core.services.networking.encrypted_proxy.hysteria2_config import build_hysteria_config
+
+# Models
+from core.models.pydantic_models.HysteriaData import HysteriaData
+from core.models.Result import Result, ResultError
+# from core.models.BaseProfile import BaseProfile
+from core.models.session.SessionProfile import SessionProfile
+from core.models.system.SystemProfile import SystemProfile
+from core.models.manage.session_management import get_session
+from core.models.manage.pydantic_management import pydantic_management
+# errors & observers
+from core.Constants import Constants
+from core.errors.logger import logger
+from core.Errors import MissingSubscriptionError
+from core.observers.ConnectionObserver import ConnectionObserver
+
+
+# generic
+from pydantic import ValidationError
+from typing import Union, Optional
+
+def configure_singbox(
+ profile: Union[SessionProfile, SystemProfile],
+ connection_observer: Optional[ConnectionObserver] = None,
+) -> Result:
+
+ ###################################
+ # PREP PAYLOADS
+ ###################################
+ protocol = profile.connection.code
+ profile_sudo_filepath = profile.get_system_config_path()
+ profile_regular_filepath = profile.get_config_path()
+ operator_id = profile.location.operator_id
+
+ logger.info(f"We're doing the protocol {protocol}, operator id of {operator_id}, and have a system path of {profile_regular_filepath}")
+
+ if not profile.has_subscription():
+ raise MissingSubscriptionError()
+
+ url = f"{Constants.SP_API_BASE_URL}/subscriptions/current/operator-proxies"
+ payload = {
+ 'operator_id': operator_id,
+ 'protocol': protocol,
+ }
+
+ ###################################
+ # SEND TO THE API
+ ###################################
+ logger.info("Sending to the API..")
+ config_results = connect.single_endpoint(
+ method="post",
+ url=url,
+ observer=connection_observer,
+ payload=payload,
+ billing_code=profile.subscription.billing_code
+ )
+
+ ###################################
+ # VERIFY THE API'S REPLY
+ ###################################
+ # this is a bad API reply, and NOT a subscription error:
+ if not config_results.valid:
+ return config_results
+
+ raw_response = config_results.data
+
+ # extract 'data' out of reply:
+ data = raw_response.get('data', raw_response)
+ logger.info(f"We got back from the API: {data}")
+
+ if not data:
+ return Result(valid=False, error_type=ResultError.INVALID_API_REPLY, error_msg="Server replied with blank or invalid data.")
+
+ ###################################
+ # VERIFY LOCATION
+ ###################################
+ location_country_code= data.get('location_country_code')
+ location_city_code= data.get('location_city_code')
+ matched_location = get_profile_location_data(
+ country_code=location_country_code,
+ city_code=location_city_code
+ )
+ if not matched_location:
+ return Result(valid=False, error_type=ResultError.INVALID_API_REPLY, error_msg="Server replied with a location that doesn't match your sync data.")
+
+ if profile.location != matched_location:
+ error_msg = f"Profile's location doesn't match. Your local data is {profile.location.country_code}_{profile.location.city_code} compared to server's {location_country_code}_{location_city_code}"
+ return Result(valid=False, error_type=ResultError.INVALID_API_REPLY, error_msg=error_msg)
+ logger.info(f"The location {profile.location.country_code} matched our local SQL")
+
+ ###################################
+ # VERIFY/SETUP/SAVE PYDANTIC MODEL
+ ###################################
+ session = get_session()
+ try:
+ if protocol == "hysteria2":
+ validated_data = HysteriaData(**data, context={"db": session})
+
+ logger.info(f"We created the model for {protocol}.")
+ except ValidationError as e:
+ for error in e.errors():
+ if error['type'] == 'missing':
+ error_msg = f"Required field '{error['loc'][0]}' is missing"
+ logger.error(error_msg)
+ else:
+ error_msg = f"Field '{error['loc'][0]}' error: {error['msg']}"
+ logger.error(error_msg)
+ return Result(valid=False, error_type=ResultError.INVALID_API_REPLY, error_msg=error_msg)
+
+ saved_raw_data = pydantic_management.save_to_regular_folder(validated_data, f"{profile_regular_filepath}/raw_setup.json")
+ logger.info(f"Saved the raw data? {saved_raw_data}")
+
+ ###################################
+ # PREP REAL CONFIG
+ ###################################
+ random_port = ports.get_random_available_port()
+
+ if protocol == "hysteria2":
+ real_config = build_hysteria_config(
+ username=validated_data.username,
+ password=validated_data.password,
+ server_host=validated_data.operator_hysteria2_host,
+ socks5_port=random_port,
+ server_ip=validated_data.server_ip
+ )
+
+
+ ###################################
+ # SAVE REAL CONFIG
+ ###################################
+ # goes in a sudo protected folder & prompts for password:
+ saved = pydantic_management.save_to_sudo_folder(real_config, f'{profile_sudo_filepath}/config.json')
+
+ if saved:
+ logger.info("Successfully saved the config.")
+ return Result(valid=True)
+ else:
+ error_msg = "Could not save the configuration."
+ return Result(valid=False, error_type=ResultError.FILE_SYSTEM, error_msg=error_msg)
diff --git a/core/services/networking/systemwide/encrypted_proxy/process_closure_tools.py b/core/services/networking/systemwide/encrypted_proxy/process_closure_tools.py
index a194c86..858b8b7 100644
--- a/core/services/networking/systemwide/encrypted_proxy/process_closure_tools.py
+++ b/core/services/networking/systemwide/encrypted_proxy/process_closure_tools.py
@@ -5,7 +5,10 @@ from core.models.Result import Result, ResultError
from core.errors.logger import logger
from core.services.networking.systemwide.encrypted_proxy import singbox
from core.utils.basic_operations import process_tools
-import subprocess
+from core.utils.run_commands import run_generic_command
+
+
+# import subprocess
from typing import Callable, cast
import time
@@ -28,7 +31,7 @@ def _try_with_permission_fallback(operation: Callable, interface: str) -> Result
try:
return operation()
except SudoScript as e:
- logger.error("The Sudo Scripts giving power to kill this are being denied permission, before we flag this, let's see if the proxy is active, which does NOT need permission to check,")
+ logger.error(f"The Sudo Scripts giving power to kill this are being denied permission, before we flag this, let's see if the proxy is active, which does NOT need permission to check, {e}")
existance = check_interface_exists(interface)
if not existance.valid and existance.error_type == ResultError.INTERFACE:
not_existing = "Permission denied, but the proxy interface is down, so this is acceptable"
@@ -97,9 +100,13 @@ def _shut_down_by_known_process_id(process_id: int) -> Result:
def get_any_pid_with_the_phrase(which_application: str) -> Result:
- result = subprocess.run(['pgrep', '-a', which_application], capture_output=True, text=True)
- if result.returncode == 0:
- lines = result.stdout.strip().split('\n')
+ # result = subprocess.run(['pgrep', '-a', which_application], capture_output=True, text=True)
+ command = ['pgrep', '-a', which_application]
+ human_readable_goal = "Find a pid by the phrase"
+ result = run_generic_command(command, human_readable_goal, timeout=7)
+
+ if result.valid:
+ lines = result.data.strip().split('\n')
pid_data = [{'pid': int(line.split()[0]), 'command': line} for line in lines if line]
if pid_data:
return Result(valid=True, data=pid_data) # List of dicts with pid and full command
@@ -111,10 +118,14 @@ def get_any_pid_with_the_phrase(which_application: str) -> Result:
# Function is public because it works for any app
def get_pid_by_exact_match(which_application: str) -> Result:
- result = subprocess.run(['pgrep', '-x', which_application], capture_output=True, text=True)
- if result.returncode == 0:
+ # result = subprocess.run(['pgrep', '-x', which_application], capture_output=True, text=True)
+ command = ['pgrep', '-x', which_application]
+ human_readable_goal = "Get a pid by an exact match"
+ result = run_generic_command(command, human_readable_goal, timeout=7)
+
+ if result.valid: # equivalent: (returncode == 0):
try:
- process_id = int(result.stdout.strip()) # Strip newline, convert to int
+ process_id = int(result.data.strip()) # Strip newline, convert to int
return Result(valid=True, data=process_id)
except ValueError:
return Result(valid=False, error_type=ResultError.INVALID_INPUT, message="Could not parse PID")
diff --git a/core/services/networking/systemwide/encrypted_proxy/singbox_runner.py b/core/services/networking/systemwide/encrypted_proxy/singbox_runner.py
index 1d93c1c..a728698 100644
--- a/core/services/networking/systemwide/encrypted_proxy/singbox_runner.py
+++ b/core/services/networking/systemwide/encrypted_proxy/singbox_runner.py
@@ -13,8 +13,8 @@ from essentials.observers.ConnectionObserver import ConnectionObserver
# generic
-import subprocess
import time
+from typing import Optional
QUANTITY_OF_ATTEMPTS = 2
@@ -56,22 +56,26 @@ def end_singbox(
connection_observer: Optional[ConnectionObserver] = None
) -> Result:
- function_name = "END_SINGBOX"
-
current_state = SystemState.get()
if not current_state:
return Result(valid=False, error_type=ResultError.MISSING_FILE, message="Already disabled or missing State JSON.")
- dead_singbox = orchestrate_closing(
- current_state=current_state,
- interface_name=Constants.SINGBOX_TUN_IF
- )
+ # kill the real singbox
+ try:
+ killed_existing = orchestrate_closing(
+ current_state=current_state,
+ interface_name=Constants.SINGBOX_TUN_IF
+ )
+ except RuntimeError as e: # Interface error raised by process_closure_tool's is_tunnel_active
+ return Result(valid=False, error_type=ResultError.INTERFACE, error_msg=str(e))
- if dead_singbox.valid:
+ # Wipe the JSON to reflect reality,
+ if killed_existing.valid:
logger.info("Successfully took down Singbox tunnel")
SystemState.dissolve()
-
-
+ return Result(valid=True)
+ else:
+ return Result(valid=False, error_type=ResultError.SINGBOX, error_msg="Could not disable singbox")
def start_singbox(
profile_id: int,
@@ -101,6 +105,11 @@ def start_singbox(
if each_requirement is None:
return Result(valid=False, error_type=ResultError.INVALID_INPUT, message=f"Invalid inputs into {function_name}")
+ # ========== KILL IT IF ALREADY UP ==========
+ killed_pre_existing = end_singbox(connection_observer)
+ if not killed_pre_existing.valid:
+ return killed_pre_existing
+
# ============= START PROCESS =============
activation_result = _attempt_start_with_retry(config_path=config_path, quantity_of_attempts=QUANTITY_OF_ATTEMPTS)
@@ -121,7 +130,7 @@ def start_singbox(
# ============= SETUP STATE =============
# Even if firewall is off, we want to save the fact we turned Singbox on, before we raise errors.
- logger.info(f"Setting State JSON with INTENDED firewall & Dns settings")
+ logger.info("Setting State JSON with INTENDED firewall & Dns settings")
current_state = SystemStateController.create(
profile_id=profile_id,
firewalled=True, # intended setting, not result yet
diff --git a/core/services/networking/systemwide/general_tools/interface_tools.py b/core/services/networking/systemwide/general_tools/interface_tools.py
index b4528e1..338b6c3 100644
--- a/core/services/networking/systemwide/general_tools/interface_tools.py
+++ b/core/services/networking/systemwide/general_tools/interface_tools.py
@@ -1,18 +1,24 @@
from core.models.Result import Result, ResultError
+from core.utils.run_commands import run_generic_command
import subprocess
import re
def check_interface_exists(interface: str) -> Result:
- result = subprocess.run(['ip', 'link', 'show', interface], capture_output=True)
- if result.returncode != 0:
+ command = ['ip', 'link', 'show', interface]
+ # result = subprocess.run(, capture_output=True)
+ human_readable_goal = "Checking if the interface exists"
+ result = run_generic_command(command, human_readable_goal, timeout=5)
+
+ if not result.valid:
return Result(valid=False, error_type=ResultError.NOT_SUPPORTED, message=f"We could not run the command to even check the interface. {result.stdout}")
- elif 'does not exist' in result.stderr.decode():
+
+ elif 'does not exist' in result.data: # .stderr.decode() was old version
return Result(valid=False, error_type=ResultError.INTERFACE, message="Interface does not exist.")
else:
# Check if interface name exists in the expected format
# Pattern: "digits: interface_name: "
- output = result.stdout.decode()
+ output = result.data #.decode()
pattern = rf'^\d+:\s+{re.escape(interface)}:\s+<[^>]+>'
if re.search(pattern, output, re.MULTILINE):
return Result(valid=True, data=output)