Compare commits
7 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3232159d29 | |||
| 47714c9622 | |||
| eba7ae831f | |||
| 74a68fcbd0 | |||
| 22b9427ca5 | |||
| ee8c0b83db | |||
| f948a410b6 |
@@ -0,0 +1,38 @@
|
|||||||
|
# SOCKS Manager - 环境变量示例
|
||||||
|
# 复制为 .env 并按需修改:cp .env.example .env
|
||||||
|
|
||||||
|
# ── 管理员账号 ──────────────────────────────────
|
||||||
|
SM_ADMIN_USER=admin
|
||||||
|
# 生产环境务必修改!首次启动如未设置会自动使用 admin123
|
||||||
|
SM_ADMIN_PASSWORD=admin123
|
||||||
|
|
||||||
|
# ── Flask 安全 ──────────────────────────────────
|
||||||
|
# 生成: python3 -c "import secrets; print(secrets.token_urlsafe(48))"
|
||||||
|
SM_SECRET_KEY=change-me-in-production
|
||||||
|
|
||||||
|
# ── Web 面板监听 ────────────────────────────────
|
||||||
|
SM_HOST=0.0.0.0
|
||||||
|
SM_PORT=5000
|
||||||
|
SM_WORKERS=1
|
||||||
|
SM_TIMEOUT=30
|
||||||
|
|
||||||
|
# ── 会话 Cookie ─────────────────────────────────
|
||||||
|
# 跨域/用域名访问时设置为你的域名或 IP
|
||||||
|
# SM_SESSION_DOMAIN=
|
||||||
|
# HTTPS 部署时设为 true
|
||||||
|
SM_COOKIE_SECURE=false
|
||||||
|
SM_SESSION_NAME=sm_session
|
||||||
|
|
||||||
|
# ── 数据库 ─────────────────────────────────────
|
||||||
|
SM_DB_URI=sqlite:///socks_manager.db
|
||||||
|
|
||||||
|
# ── 备份 ───────────────────────────────────────
|
||||||
|
SM_BACKUP_DIR=./backups/
|
||||||
|
|
||||||
|
# ── 日志 ───────────────────────────────────────
|
||||||
|
SM_LOG_LEVEL=INFO
|
||||||
|
|
||||||
|
# ── SOCKS5 引擎 ────────────────────────────────
|
||||||
|
# 是否启用后台 sync 线程(探活死掉的 SOCKS5 实例)
|
||||||
|
SM_SYNC_LOOP_ENABLE=true
|
||||||
|
SM_SYNC_INTERVAL=30
|
||||||
@@ -187,66 +187,53 @@ Flask 内置开发服务器不适合生产环境。使用 **Gunicorn**(生产
|
|||||||
# 安装 gunicorn
|
# 安装 gunicorn
|
||||||
pip install gunicorn --break-system-packages
|
pip install gunicorn --break-system-packages
|
||||||
|
|
||||||
# 启动(4 worker 进程,适合 4 核服务器)
|
# 最小启动(不推荐生产,worker 数量 = CPU 核心数,单核机器用 -w 1)
|
||||||
gunicorn -w 4 -b 0.0.0.0:5000 --timeout 120 run:app
|
gunicorn -w 1 -b 0.0.0.0:5000 run:app
|
||||||
|
|
||||||
# 推荐参数(更长超时 + 优雅关闭)
|
# 推荐参数:30s 超时 + 优雅关闭 + 周期性回收 worker 防内存泄漏
|
||||||
gunicorn -w 2 -b 0.0.0.0:5000 \
|
gunicorn -w 1 -b 0.0.0.0:5000 \
|
||||||
--timeout 300 \
|
--timeout 30 \
|
||||||
--graceful-timeout 30 \
|
--graceful-timeout 30 \
|
||||||
--max-requests 10000 \
|
--max-requests 1000 \
|
||||||
--max-requests-jitter 1000 \
|
--max-requests-jitter 100 \
|
||||||
|
--limit-request-line 8190 \
|
||||||
|
--limit-request-fields 100 \
|
||||||
run:app
|
run:app
|
||||||
```
|
```
|
||||||
|
|
||||||
> **注意**:SOCKS5 代理实例运行在独立的 asyncio 线程中,不受 Gunicorn worker 数量影响。
|
> **注意**:SOCKS5 代理实例运行在独立的 asyncio 线程中,不受 Gunicorn worker 数量影响。
|
||||||
|
> **生产请用 systemd 托管**,见下节;裸 gunicorn 没有开机自启和崩溃重启。
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|
||||||
### 2. Systemd 服务配置(推荐)
|
### 2. Systemd 服务配置(推荐)
|
||||||
|
|
||||||
创建 `/etc/systemd/system/socks-manager.service`:
|
服务单元文件的**唯一权威源**在仓库 `deploy/socks-manager.service`。
|
||||||
|
所有参数、调优注释、`MemoryMax` 限制都集中在那一份,生产机只放一份拷贝。
|
||||||
|
|
||||||
```ini
|
部署步骤:
|
||||||
[Unit]
|
|
||||||
Description=SOCKS5 Proxy Manager
|
|
||||||
After=network.target
|
|
||||||
|
|
||||||
[Service]
|
```bash
|
||||||
Type=notify
|
# 1. 在本地修改 deploy/socks-manager.service
|
||||||
User=root
|
# 2. git commit + push
|
||||||
WorkingDirectory=/opt/socks-manager
|
# 3. 在生产机 (/opt/socks-manager) 拉取并执行:
|
||||||
Environment="PATH=/usr/local/bin:/usr/bin:/bin"
|
git pull
|
||||||
|
sudo bash deploy/install.sh
|
||||||
# 生产环境变量(按需修改)
|
|
||||||
Environment="SM_ADMIN_PASSWORD=你的强密码"
|
|
||||||
Environment="SM_SECRET_KEY=生产环境用openssl rand -hex 32生成"
|
|
||||||
Environment="SM_LOG_LEVEL=WARN"
|
|
||||||
# Environment="SM_COOKIE_SECURE=true" # HTTPS 时开启
|
|
||||||
# Environment="SM_SESSION_DOMAIN=你的域名或IP" # 跨域时设置
|
|
||||||
|
|
||||||
ExecStart=/usr/local/bin/gunicorn -w 2 -b 0.0.0.0:5000 \
|
|
||||||
--timeout 300 \
|
|
||||||
--graceful-timeout 30 \
|
|
||||||
run:app
|
|
||||||
|
|
||||||
# 优雅关闭:先给 SOCKS5 实例时间关闭活跃连接
|
|
||||||
ExecStop=/bin/kill -s TERM $MAINPID
|
|
||||||
TimeoutStopSec=60
|
|
||||||
|
|
||||||
Restart=always
|
|
||||||
RestartSec=5
|
|
||||||
|
|
||||||
[Install]
|
|
||||||
WantedBy=multi-user.target
|
|
||||||
```
|
```
|
||||||
|
|
||||||
启动服务:
|
`install.sh` 做的事:复制 unit 到 `/etc/systemd/system/` → `daemon-reload` → `enable` (开机自启) → `restart` → 等 15s 验证 active。
|
||||||
|
|
||||||
|
完整 unit 见 [`deploy/socks-manager.service`](deploy/socks-manager.service),关键参数说明:
|
||||||
|
- `-w 1` / `--timeout 30` / `--graceful-timeout 30` — 1 核 + 1GB 内存机型的基线,按 `nproc` 调整 `-w`
|
||||||
|
- `--max-requests 1000` + jitter — 周期性回收 worker 防内存泄漏
|
||||||
|
- `MemoryMax=512M` — cgroup 硬上限,防整台 VPS 被拖死
|
||||||
|
- `LimitNOFILE=65535` — SOCKS5 大并发客户端连接需要
|
||||||
|
|
||||||
|
常用命令:
|
||||||
```bash
|
```bash
|
||||||
systemctl daemon-reload
|
|
||||||
systemctl enable --now socks-manager
|
|
||||||
systemctl status socks-manager
|
systemctl status socks-manager
|
||||||
journalctl -u socks-manager -f # 查看实时日志
|
journalctl -u socks-manager -f # 实时日志
|
||||||
|
systemctl restart socks-manager # 改完代码后重启
|
||||||
```
|
```
|
||||||
|
|
||||||
---
|
---
|
||||||
|
|||||||
@@ -1,7 +1,20 @@
|
|||||||
"""Flask 应用工厂。"""
|
"""Flask 应用工厂。"""
|
||||||
|
import atexit
|
||||||
import os
|
import os
|
||||||
import logging
|
import logging
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
from flask import Flask
|
from flask import Flask
|
||||||
|
|
||||||
|
# 可选: 加载 .env 文件(生产部署推荐用 systemd EnvironmentFile)
|
||||||
|
try:
|
||||||
|
from dotenv import load_dotenv
|
||||||
|
_env_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".env")
|
||||||
|
if os.path.isfile(_env_path):
|
||||||
|
load_dotenv(_env_path)
|
||||||
|
except ImportError:
|
||||||
|
pass
|
||||||
|
|
||||||
from config import Config
|
from config import Config
|
||||||
from database import db
|
from database import db
|
||||||
from models import Instance # 触发表创建
|
from models import Instance # 触发表创建
|
||||||
@@ -11,6 +24,58 @@ from engine.instances import InstanceManager
|
|||||||
from services.user_service import UserService
|
from services.user_service import UserService
|
||||||
from services.backup_service import set_backup_dir
|
from services.backup_service import set_backup_dir
|
||||||
|
|
||||||
|
# ── SOCKS5 实例自动 sync 调度 ────────────────────────────────
|
||||||
|
# 只有 master 进程会跑 (gunicorn -w 1 时 = worker 0)。
|
||||||
|
# 周期调用 instance_manager.sync_instances(), 它会调 Socks5Server.is_alive()
|
||||||
|
# 探测死掉的实例并用 DB 配置自动重启, 同时修正 DB.running 字段。
|
||||||
|
#
|
||||||
|
# 之前 sync_instances 只能由 API 触发 (/api/instances/sync), 没人调就没人探活,
|
||||||
|
# "实例莫名停止" 实际是 SOCKS5 线程死掉但没人发现。
|
||||||
|
#
|
||||||
|
# 30s 周期足够短, 用户感觉不到抖动; 也不会跟 gunicorn worker timeout (30s)
|
||||||
|
# 竞争——gunicorn 在 HTTP 请求上下文里 timeout, sync_instances 在主线程跑
|
||||||
|
# 没有 gunicorn 中间件, 不会触发 worker timeout。
|
||||||
|
_sync_thread = None
|
||||||
|
_sync_lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
|
def _start_sync_loop(app):
|
||||||
|
"""启动后台 sync 线程。线程级单例, 多次调用只启一次。
|
||||||
|
|
||||||
|
gunicorn -w N 时, 每个 worker 都会调 create_app() 一次, 但实际 gunicorn 的
|
||||||
|
post_fork hook 只在 worker 进程里跑。如果用 N>1, 需要用 file lock 或
|
||||||
|
gunicorn master/worker 区分; 简单起见, 用环境变量 SM_SYNC_LOOP_ENABLE
|
||||||
|
显式开 (生产 unit 不设, 所以默认是关, 留 gunicorn 进程数=1 的安全 case)。
|
||||||
|
"""
|
||||||
|
global _sync_thread
|
||||||
|
with _sync_lock:
|
||||||
|
if _sync_thread is not None and _sync_thread.is_alive():
|
||||||
|
return
|
||||||
|
if os.environ.get("SM_SYNC_LOOP_ENABLE", "true").lower() != "true":
|
||||||
|
log.info("sync loop disabled (SM_SYNC_LOOP_ENABLE != true)")
|
||||||
|
return
|
||||||
|
interval = int(os.environ.get("SM_SYNC_INTERVAL", "30"))
|
||||||
|
if interval < 5:
|
||||||
|
interval = 5 # 5s 下限, 防止在 hot loop 里被滥用
|
||||||
|
im = app.instance_manager
|
||||||
|
|
||||||
|
def _loop():
|
||||||
|
log.info("SOCKS5 sync loop started, interval=%ds", interval)
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
im.sync_instances()
|
||||||
|
except Exception as e:
|
||||||
|
log.exception("sync_instances failed: %s", e)
|
||||||
|
time.sleep(interval)
|
||||||
|
|
||||||
|
t = threading.Thread(target=_loop, daemon=True, name="socks5-sync")
|
||||||
|
t.start()
|
||||||
|
_sync_thread = t
|
||||||
|
atexit.register(lambda: log.info("socks5-sync loop exiting (atexit)"))
|
||||||
|
|
||||||
|
|
||||||
|
log = logging.getLogger("socks.app")
|
||||||
|
|
||||||
|
|
||||||
def create_app():
|
def create_app():
|
||||||
# 确保密码已设置(首次启动自动设置默认密码)
|
# 确保密码已设置(首次启动自动设置默认密码)
|
||||||
@@ -63,4 +128,13 @@ def create_app():
|
|||||||
app.register_blueprint(web_bp)
|
app.register_blueprint(web_bp)
|
||||||
app.register_blueprint(api_bp, url_prefix="/api")
|
app.register_blueprint(api_bp, url_prefix="/api")
|
||||||
|
|
||||||
|
# 启动后台 sync 线程 (在主线程, 不在 worker 上下文)
|
||||||
|
# 必须在 register_blueprint 之后, 调一次 sync_instances 把当前已配的
|
||||||
|
# SOCKS5 实例拉起来, 再进 loop 周期探活
|
||||||
|
_start_sync_loop(app)
|
||||||
|
try:
|
||||||
|
instance_manager.sync_instances()
|
||||||
|
except Exception as e:
|
||||||
|
log.exception("initial sync_instances failed: %s", e)
|
||||||
|
|
||||||
return app
|
return app
|
||||||
|
|||||||
@@ -0,0 +1,297 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# =============================================================================
|
||||||
|
# SOCKS Manager - 一键部署脚本
|
||||||
|
# =============================================================================
|
||||||
|
# 功能:
|
||||||
|
# 1. 安装系统依赖 (python3, pip, gcc, libffi, etc.)
|
||||||
|
# 2. 创建 Python 虚拟环境并安装 pip 依赖
|
||||||
|
# 3. 生成 .env 配置(随机 SECRET_KEY、管理员密码)
|
||||||
|
# 4. 初始化数据库 & 后台线程
|
||||||
|
# 5. 生成 gunicorn + systemd 服务文件并启动
|
||||||
|
# 6. 配置防火墙(firewalld / ufw 自动识别)
|
||||||
|
#
|
||||||
|
# 用法:
|
||||||
|
# bash deploy.sh # 全自动部署 (默认 /opt/socks-manager)
|
||||||
|
# INSTALL_DIR=/opt/sm bash deploy.sh
|
||||||
|
# bash deploy.sh --no-start # 不启动服务
|
||||||
|
# ADMIN_PASS=mypassword bash deploy.sh
|
||||||
|
# =============================================================================
|
||||||
|
set -euo pipefail
|
||||||
|
|
||||||
|
# ---- 配置 -------------------------------------------------------------------
|
||||||
|
INSTALL_DIR="${INSTALL_DIR:-/opt/socks-manager}"
|
||||||
|
SERVICE_NAME="${SERVICE_NAME:-socks-manager}"
|
||||||
|
SM_HOST="${SM_HOST:-0.0.0.0}"
|
||||||
|
SM_PORT="${SM_PORT:-5000}"
|
||||||
|
SM_WORKERS="${SM_WORKERS:-1}"
|
||||||
|
ADMIN_USER="${ADMIN_USER:-admin}"
|
||||||
|
ADMIN_PASS="${ADMIN_PASS:-admin123}"
|
||||||
|
|
||||||
|
START_SERVICE=true
|
||||||
|
SKIP_EXISTING=false
|
||||||
|
|
||||||
|
for arg in "$@"; do
|
||||||
|
case "$arg" in
|
||||||
|
--no-start) START_SERVICE=false ;;
|
||||||
|
--skip-existing) SKIP_EXISTING=true ;;
|
||||||
|
-h|--help)
|
||||||
|
sed -n '2,25p' "$0"
|
||||||
|
exit 0
|
||||||
|
;;
|
||||||
|
*)
|
||||||
|
echo "未知参数: $arg" >&2; exit 1 ;;
|
||||||
|
esac
|
||||||
|
done
|
||||||
|
|
||||||
|
# ---- 工具函数 ---------------------------------------------------------------
|
||||||
|
RED='\033[0;31m'; GREEN='\033[0;32m'; YELLOW='\033[1;33m'; NC='\033[0m'
|
||||||
|
info() { echo -e "${GREEN}[INFO]${NC} $*"; }
|
||||||
|
warn() { echo -e "${YELLOW}[WARN]${NC} $*"; }
|
||||||
|
error() { echo -e "${RED}[ERROR]${NC} $*" >&2; }
|
||||||
|
|
||||||
|
must_run_as_root() {
|
||||||
|
if [[ $EUID -ne 0 ]]; then
|
||||||
|
error "此脚本必须以 root 运行 (sudo $0)"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
|
||||||
|
detect_os() {
|
||||||
|
if [[ -f /etc/os-release ]]; then
|
||||||
|
. /etc/os-release
|
||||||
|
OS_ID="$ID"
|
||||||
|
OS_VER="$VERSION_ID"
|
||||||
|
elif command -v apt-get &>/dev/null; then
|
||||||
|
OS_ID="debian"
|
||||||
|
elif command -v yum &>/dev/null; then
|
||||||
|
OS_ID="centos"
|
||||||
|
else
|
||||||
|
error "无法识别操作系统"
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
info "操作系统: $OS_ID $OS_VER"
|
||||||
|
}
|
||||||
|
|
||||||
|
pkg_install() {
|
||||||
|
info "安装系统包: $*"
|
||||||
|
case "$OS_ID" in
|
||||||
|
debian|ubuntu)
|
||||||
|
export DEBIAN_FRONTEND=noninteractive
|
||||||
|
apt-get update -qq
|
||||||
|
apt-get install -y -qq "$@"
|
||||||
|
;;
|
||||||
|
centos|rhel|rocky|almalinux|ol)
|
||||||
|
yum install -y -q "$@"
|
||||||
|
;;
|
||||||
|
*)
|
||||||
|
error "不支持的包管理器: $OS_ID"
|
||||||
|
exit 1
|
||||||
|
;;
|
||||||
|
esac
|
||||||
|
}
|
||||||
|
|
||||||
|
# ---- 1. 前置检查 ------------------------------------------------------------
|
||||||
|
must_run_as_root
|
||||||
|
detect_os
|
||||||
|
|
||||||
|
info "安装目录: $INSTALL_DIR"
|
||||||
|
info "服务名: $SERVICE_NAME"
|
||||||
|
info "监听: $SM_HOST:$SM_PORT"
|
||||||
|
info "管理员: $ADMIN_USER"
|
||||||
|
|
||||||
|
# ---- 2. 安装系统依赖 --------------------------------------------------------
|
||||||
|
SYS_PACKAGES=(
|
||||||
|
python3 python3-pip python3-venv python3-dev
|
||||||
|
gcc libffi-dev libssl-dev
|
||||||
|
curl
|
||||||
|
)
|
||||||
|
|
||||||
|
pkg_install "${SYS_PACKAGES[@]}"
|
||||||
|
|
||||||
|
PYTHON_BIN="$(command -v python3)"
|
||||||
|
PIP_BIN="$(command -v pip3 || command -v pip)"
|
||||||
|
info "Python: $PYTHON_BIN"
|
||||||
|
info "Pip: $PIP_BIN"
|
||||||
|
|
||||||
|
# ---- 3. 部署项目文件 --------------------------------------------------------
|
||||||
|
SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)"
|
||||||
|
|
||||||
|
if [[ "$SCRIPT_DIR" != "$INSTALL_DIR" ]]; then
|
||||||
|
if [[ -d "$INSTALL_DIR" ]] && $SKIP_EXISTING; then
|
||||||
|
warn "安装目录已存在且 --skip-existing,跳过文件复制"
|
||||||
|
else
|
||||||
|
info "复制项目文件到 $INSTALL_DIR"
|
||||||
|
mkdir -p "$INSTALL_DIR"
|
||||||
|
rsync -a --delete \
|
||||||
|
--exclude='.git/' \
|
||||||
|
--exclude='__pycache__/' \
|
||||||
|
--exclude='*.pyc' \
|
||||||
|
--exclude='*.db' \
|
||||||
|
--exclude='backups/' \
|
||||||
|
--exclude='.env' \
|
||||||
|
--exclude='gunicorn.pid' \
|
||||||
|
"$SCRIPT_DIR/" "$INSTALL_DIR/"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
cd "$INSTALL_DIR"
|
||||||
|
|
||||||
|
# 确保运行时目录
|
||||||
|
mkdir -p backups
|
||||||
|
chmod 700 backups
|
||||||
|
|
||||||
|
# ---- 4. 虚拟环境 + pip 依赖 ------------------------------------------------
|
||||||
|
VENV_DIR="$INSTALL_DIR/.venv"
|
||||||
|
|
||||||
|
if [[ -d "$VENV_DIR" ]] && $SKIP_EXISTING; then
|
||||||
|
warn "虚拟环境已存在且 --skip-existing,跳过依赖安装"
|
||||||
|
else
|
||||||
|
info "创建 Python 虚拟环境: $VENV_DIR"
|
||||||
|
$PYTHON_BIN -m venv "$VENV_DIR"
|
||||||
|
|
||||||
|
info "升级 pip + 安装 Python 依赖"
|
||||||
|
"$VENV_DIR/bin/pip" install --upgrade pip setuptools wheel
|
||||||
|
"$VENV_DIR/bin/pip" install -r requirements.txt
|
||||||
|
fi
|
||||||
|
|
||||||
|
PYTHON="$VENV_DIR/bin/python"
|
||||||
|
GUNICORN="$VENV_DIR/bin/gunicorn"
|
||||||
|
|
||||||
|
# ---- 5. 生成 .env 配置 ------------------------------------------------------
|
||||||
|
ENV_FILE="$INSTALL_DIR/.env"
|
||||||
|
|
||||||
|
if [[ ! -f "$ENV_FILE" ]] || ! $SKIP_EXISTING; then
|
||||||
|
info "生成 .env 配置文件"
|
||||||
|
SECRET_KEY="$($PYTHON -c "import secrets; print(secrets.token_urlsafe(48))")"
|
||||||
|
cat > "$ENV_FILE" <<EOF
|
||||||
|
# SOCKS Manager - 自动生成于 $(date -Iseconds)
|
||||||
|
# 所有变量可参考 .env.example
|
||||||
|
SM_ADMIN_USER=$ADMIN_USER
|
||||||
|
SM_ADMIN_PASSWORD=$ADMIN_PASS
|
||||||
|
SM_SECRET_KEY=$SECRET_KEY
|
||||||
|
SM_HOST=$SM_HOST
|
||||||
|
SM_PORT=$SM_PORT
|
||||||
|
SM_WORKERS=$SM_WORKERS
|
||||||
|
SM_TIMEOUT=30
|
||||||
|
SM_LOG_LEVEL=INFO
|
||||||
|
SM_SYNC_LOOP_ENABLE=true
|
||||||
|
SM_SYNC_INTERVAL=30
|
||||||
|
EOF
|
||||||
|
chmod 600 "$ENV_FILE"
|
||||||
|
info "SECRET_KEY 已随机生成,管理员密码: $ADMIN_PASS"
|
||||||
|
fi
|
||||||
|
|
||||||
|
# ---- 6. 初始化数据库 --------------------------------------------------------
|
||||||
|
info "初始化数据库"
|
||||||
|
$PYTHON -c "
|
||||||
|
import os, sys
|
||||||
|
os.environ['SM_ADMIN_PASSWORD'] = '$ADMIN_PASS'
|
||||||
|
os.environ['SM_ADMIN_USER'] = '$ADMIN_USER'
|
||||||
|
sys.path.insert(0, '$INSTALL_DIR')
|
||||||
|
from app import create_app
|
||||||
|
app = create_app()
|
||||||
|
print('[init] database ready, instances synced')
|
||||||
|
"
|
||||||
|
|
||||||
|
# ---- 7. 生成 systemd 服务文件 ----------------------------------------------
|
||||||
|
SERVICE_FILE="/etc/systemd/system/${SERVICE_NAME}.service"
|
||||||
|
|
||||||
|
info "生成 systemd 服务文件: $SERVICE_FILE"
|
||||||
|
cat > "$SERVICE_FILE" <<EOF
|
||||||
|
[Unit]
|
||||||
|
Description=SOCKS5 Proxy Manager
|
||||||
|
After=network.target
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=simple
|
||||||
|
User=root
|
||||||
|
Group=root
|
||||||
|
WorkingDirectory=$INSTALL_DIR
|
||||||
|
EnvironmentFile=$ENV_FILE
|
||||||
|
ExecStart=$GUNICORN -c gunicorn_config.py run:app
|
||||||
|
ExecReload=/bin/kill -s HUP \$MAINPID
|
||||||
|
ExecStop=/bin/kill -s TERM \$MAINPID
|
||||||
|
Restart=always
|
||||||
|
RestartSec=5
|
||||||
|
TimeoutStopSec=60
|
||||||
|
# 资源限制
|
||||||
|
MemoryMax=512M
|
||||||
|
MemoryHigh=384M
|
||||||
|
LimitNOFILE=65535
|
||||||
|
# 安全加固
|
||||||
|
NoNewPrivileges=yes
|
||||||
|
ProtectSystem=full
|
||||||
|
ProtectHome=true
|
||||||
|
ReadWritePaths=$INSTALL_DIR
|
||||||
|
PrivateTmp=yes
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=multi-user.target
|
||||||
|
EOF
|
||||||
|
|
||||||
|
systemctl daemon-reload
|
||||||
|
systemctl enable "$SERVICE_NAME"
|
||||||
|
info "systemd 服务已注册并设为开机自启"
|
||||||
|
|
||||||
|
# ---- 8. 防火墙配置 ---------------------------------------------------------
|
||||||
|
_configure_firewall() {
|
||||||
|
local port="$1"
|
||||||
|
if command -v ufw &>/dev/null && ufw status | grep -q "Status: active" 2>/dev/null; then
|
||||||
|
if ! ufw status | grep -q "${port}/tcp"; then
|
||||||
|
ufw allow "${port}/tcp" comment "socks-manager web" >/dev/null
|
||||||
|
info "ufw 已放行 ${port}/tcp (Web 面板)"
|
||||||
|
fi
|
||||||
|
elif command -v firewall-cmd &>/dev/null && firewall-cmd --state &>/dev/null; then
|
||||||
|
if ! firewall-cmd --list-ports | grep -q "${port}/tcp"; then
|
||||||
|
firewall-cmd --permanent --add-port="${port}/tcp" >/dev/null
|
||||||
|
firewall-cmd --reload >/dev/null
|
||||||
|
info "firewalld 已放行 ${port}/tcp (Web 面板)"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
}
|
||||||
|
_configure_firewall "$SM_PORT"
|
||||||
|
|
||||||
|
# ---- 9. 启动服务 -----------------------------------------------------------
|
||||||
|
if $START_SERVICE; then
|
||||||
|
info "启动 $SERVICE_NAME 服务"
|
||||||
|
systemctl restart "$SERVICE_NAME"
|
||||||
|
sleep 3
|
||||||
|
|
||||||
|
if systemctl is-active --quiet "$SERVICE_NAME"; then
|
||||||
|
info "$SERVICE_NAME 服务运行正常 ✓"
|
||||||
|
else
|
||||||
|
error "$SERVICE_NAME 服务启动失败!"
|
||||||
|
echo "---------- journalctl 最后 40 行 ----------"
|
||||||
|
journalctl -u "$SERVICE_NAME" --no-pager -n 40
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# 健康检查
|
||||||
|
if curl -sf -o /dev/null "http://127.0.0.1:${SM_PORT}/login"; then
|
||||||
|
info "HTTP 健康检查通过 ✓ (http://127.0.0.1:${SM_PORT}/login)"
|
||||||
|
else
|
||||||
|
warn "HTTP 健康检查失败,请检查防火墙或端口配置"
|
||||||
|
fi
|
||||||
|
fi
|
||||||
|
|
||||||
|
# ---- 10. 输出汇总 -----------------------------------------------------------
|
||||||
|
echo ""
|
||||||
|
echo "================================================================"
|
||||||
|
echo " SOCKS Manager 部署完成!"
|
||||||
|
echo "================================================================"
|
||||||
|
echo " Web UI: http://<服务器IP>:${SM_PORT}"
|
||||||
|
echo " 账号: ${ADMIN_USER} / ${ADMIN_PASS}"
|
||||||
|
echo " 安装目录: ${INSTALL_DIR}"
|
||||||
|
echo " 服务名: ${SERVICE_NAME}"
|
||||||
|
echo " 配置文件: ${ENV_FILE}"
|
||||||
|
echo ""
|
||||||
|
echo " 常用命令:"
|
||||||
|
echo " systemctl status $SERVICE_NAME # 查看状态"
|
||||||
|
echo " systemctl restart $SERVICE_NAME # 重启服务"
|
||||||
|
echo " journalctl -u $SERVICE_NAME -f # 实时日志"
|
||||||
|
echo ""
|
||||||
|
echo " 提示:"
|
||||||
|
echo " - SOCKS5 代理实例在 Web 面板中创建和管理"
|
||||||
|
echo " - 代理端口需要手动在防火墙放行(面板会提示)"
|
||||||
|
echo " - 登录后请立即修改默认管理员密码"
|
||||||
|
echo "================================================================"
|
||||||
Executable
+43
@@ -0,0 +1,43 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
# 把 deploy/socks-manager.service 同步到生产机并重启。
|
||||||
|
#
|
||||||
|
# 使用: 在生产机 /opt/socks-manager 下执行
|
||||||
|
# bash deploy/install.sh
|
||||||
|
#
|
||||||
|
# 假定: 仓库已 git pull 到最新,当前在 /opt/socks-manager 工作目录下。
|
||||||
|
|
||||||
|
set -euo pipefail
|
||||||
|
|
||||||
|
SERVICE_FILE="deploy/socks-manager.service"
|
||||||
|
DEST="/etc/systemd/system/socks-manager.service"
|
||||||
|
|
||||||
|
if [[ ! -f "$SERVICE_FILE" ]]; then
|
||||||
|
echo "ERROR: $SERVICE_FILE not found. Run from /opt/socks-manager." >&2
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
echo "[1/5] install unit -> $DEST"
|
||||||
|
install -m 0644 "$SERVICE_FILE" "$DEST"
|
||||||
|
|
||||||
|
echo "[2/5] daemon-reload"
|
||||||
|
systemctl daemon-reload
|
||||||
|
|
||||||
|
echo "[3/5] enable (开机自启)"
|
||||||
|
systemctl enable socks-manager.service
|
||||||
|
|
||||||
|
echo "[4/5] restart"
|
||||||
|
systemctl restart socks-manager.service
|
||||||
|
|
||||||
|
echo "[5/5] wait for active (max 15s)"
|
||||||
|
for i in $(seq 1 15); do
|
||||||
|
if systemctl is-active --quiet socks-manager.service; then
|
||||||
|
echo "active after ${i}s"
|
||||||
|
systemctl status socks-manager.service --no-pager | head -15
|
||||||
|
exit 0
|
||||||
|
fi
|
||||||
|
sleep 1
|
||||||
|
done
|
||||||
|
|
||||||
|
echo "ERROR: service did not become active in 15s" >&2
|
||||||
|
journalctl -u socks-manager.service -n 80 --no-pager >&2
|
||||||
|
exit 1
|
||||||
@@ -0,0 +1,63 @@
|
|||||||
|
# SOCKS Manager — systemd unit (authoritative source)
|
||||||
|
#
|
||||||
|
# 这个文件是服务器上 /etc/systemd/system/socks-manager.service 的唯一定义源。
|
||||||
|
# 修改流程:
|
||||||
|
# 1. 编辑本文件
|
||||||
|
# 2. 提交并 push 到 origin
|
||||||
|
# 3. 在生产机执行: bash deploy/install.sh
|
||||||
|
#
|
||||||
|
# 调优说明(按机器规格):
|
||||||
|
# -w N gunicorn worker 数。1 核 CPU 建议 1,2 核建议 2。
|
||||||
|
# --timeout 30 HTTP 请求行/头/体的总超时(秒)。网络扫描器发 PRI * HTTP/2.0 这类
|
||||||
|
# 半截请求会占满 worker,30s 既能容忍慢客户端也能及时回收。
|
||||||
|
# --graceful-timeout 30 systemd 停止时给 worker 写完响应的最大等待时间。
|
||||||
|
# --max-requests 1000 worker 处理 N 个请求后自杀重启,防止内存泄漏累积。
|
||||||
|
# --max-requests-jitter 100 错峰重启,避免所有 worker 同时重启。
|
||||||
|
# MemoryMax=512M cgroup 内存硬上限,超限立即 OOM kill,保护宿主机。
|
||||||
|
# LimitNOFILE=65535 SOCKS5 代理每个实例一个监听 socket + 大量客户端连接,
|
||||||
|
# 默认 1024 远远不够。
|
||||||
|
|
||||||
|
[Unit]
|
||||||
|
Description=SOCKS5 Proxy Manager
|
||||||
|
After=network.target
|
||||||
|
|
||||||
|
[Service]
|
||||||
|
Type=notify
|
||||||
|
User=root
|
||||||
|
WorkingDirectory=/opt/socks-manager
|
||||||
|
Environment="PATH=/usr/local/bin:/usr/bin:/bin"
|
||||||
|
|
||||||
|
# 生产环境变量(按需修改)— 敏感信息建议放到 /opt/socks-manager/.env 或
|
||||||
|
# /etc/socks-manager/socks-manager.env,用 EnvironmentFile= 引入,不要写进 git。
|
||||||
|
Environment="SM_ADMIN_PASSWORD=test123"
|
||||||
|
Environment="SM_SECRET_KEY=xxxx111wwwddadljkkjkjnmooihjklhhgk"
|
||||||
|
Environment="SM_LOG_LEVEL=WARN"
|
||||||
|
# Environment="SM_COOKIE_SECURE=true" # HTTPS 时开启
|
||||||
|
# Environment="SM_SESSION_DOMAIN=你的域名或IP" # 跨域时设置
|
||||||
|
|
||||||
|
ExecStart=/usr/local/bin/gunicorn \
|
||||||
|
-w 1 \
|
||||||
|
-b 0.0.0.0:5000 \
|
||||||
|
--timeout 30 \
|
||||||
|
--graceful-timeout 30 \
|
||||||
|
--max-requests 1000 \
|
||||||
|
--max-requests-jitter 100 \
|
||||||
|
--limit-request-line 8190 \
|
||||||
|
--limit-request-fields 100 \
|
||||||
|
--limit-request-field_size 8190 \
|
||||||
|
run:app
|
||||||
|
|
||||||
|
# 优雅关闭:先给 SOCKS5 实例时间关闭活跃连接
|
||||||
|
ExecStop=/bin/kill -s TERM $MAINPID
|
||||||
|
TimeoutStopSec=60
|
||||||
|
|
||||||
|
Restart=always
|
||||||
|
RestartSec=5
|
||||||
|
|
||||||
|
# 资源限制:1GB 内存的 VPS 单 worker 留 512M 足够,超了立刻被 OOM kill 而非慢慢拖死整机
|
||||||
|
MemoryMax=512M
|
||||||
|
MemoryHigh=384M
|
||||||
|
LimitNOFILE=65535
|
||||||
|
|
||||||
|
[Install]
|
||||||
|
WantedBy=multi-user.target
|
||||||
+60
-6
@@ -150,6 +150,32 @@ class Socks5Server:
|
|||||||
with self._lock:
|
with self._lock:
|
||||||
return self._active_connections
|
return self._active_connections
|
||||||
|
|
||||||
|
def is_alive(self):
|
||||||
|
"""真实健康检查: 线程活 + server socket 在服务。
|
||||||
|
|
||||||
|
self.running 是乐观标记, start() 之后一直为 True 直到 stop() 被显式调用。
|
||||||
|
如果子线程 event loop 崩了 (例如 _record_stats 抛 TypeError 把 worker 卡死,
|
||||||
|
上游 gunicorn timeout kill 进程, 或端口被外部抢占), self.running 还是
|
||||||
|
True 但 server 实际不再 accept。
|
||||||
|
|
||||||
|
调用方应据此重启实例, 否则 DB 里 running=true 是谎言。
|
||||||
|
|
||||||
|
判定:
|
||||||
|
- thread 死了 -> 死 (无法恢复)
|
||||||
|
- server is None -> 死 (还没起来)
|
||||||
|
- server is_serving() == False -> 死 (socket 已关, asyncio.Server
|
||||||
|
不能复用同一个对象重新 serve, 必须新 Socks5Server)
|
||||||
|
"""
|
||||||
|
if not self.running:
|
||||||
|
return False
|
||||||
|
if self.thread is None or not self.thread.is_alive():
|
||||||
|
return False
|
||||||
|
if self.server is None:
|
||||||
|
return False
|
||||||
|
if not self.server.is_serving():
|
||||||
|
return False
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
import time # noqa: E402
|
import time # noqa: E402
|
||||||
|
|
||||||
@@ -185,7 +211,12 @@ class InstanceManager:
|
|||||||
self._config_cache = new_configs
|
self._config_cache = new_configs
|
||||||
|
|
||||||
def sync_instances(self):
|
def sync_instances(self):
|
||||||
"""根据数据库配置同步实例运行状态。"""
|
"""根据数据库配置同步实例运行状态。
|
||||||
|
|
||||||
|
健康检查: 对每个 _instances[name], 调用 is_alive() 验证线程 + server
|
||||||
|
socket 都活着。如果死了, 从 _instances 删掉, 用同一个 config 重新启动。
|
||||||
|
避免 40080 那种'DB 写 running=true 但端口没人 listen'的谎言状态。
|
||||||
|
"""
|
||||||
from models import Instance
|
from models import Instance
|
||||||
from database import db
|
from database import db
|
||||||
self.reload_config()
|
self.reload_config()
|
||||||
@@ -193,7 +224,18 @@ class InstanceManager:
|
|||||||
current_names = set(self._config_cache.keys())
|
current_names = set(self._config_cache.keys())
|
||||||
running_names = set(self._instances.keys())
|
running_names = set(self._instances.keys())
|
||||||
|
|
||||||
# 启动缺失的实例
|
# 健康检查: 把死的从 _instances 摘掉, 让下面的'启动缺失'逻辑接管。
|
||||||
|
# _config_cache 是从 DB reload 的权威配置, 不能动它, 否则重启循环拿不到 cfg。
|
||||||
|
dead = []
|
||||||
|
for name, srv in self._instances.items():
|
||||||
|
if not srv.is_alive():
|
||||||
|
log.warning("[%s] SOCKS5 进程不健康 (thread/socket dead), 标记待重启", name)
|
||||||
|
dead.append(name)
|
||||||
|
for name in dead:
|
||||||
|
self._instances.pop(name, None)
|
||||||
|
running_names = set(self._instances.keys())
|
||||||
|
|
||||||
|
# 启动缺失的实例 (含刚被健康检查踢出来的)
|
||||||
for name in current_names - running_names:
|
for name in current_names - running_names:
|
||||||
cfg = self._config_cache[name]
|
cfg = self._config_cache[name]
|
||||||
srv = Socks5Server(cfg, self.user_service, self.flask_app)
|
srv = Socks5Server(cfg, self.user_service, self.flask_app)
|
||||||
@@ -210,11 +252,11 @@ class InstanceManager:
|
|||||||
srv = self._instances.pop(name)
|
srv = self._instances.pop(name)
|
||||||
srv.stop()
|
srv.stop()
|
||||||
|
|
||||||
# 更新实例运行状态到数据库
|
# 更新实例运行状态到数据库 (用 is_alive() 而不是乐观 running 标记)
|
||||||
with db.app.app_context():
|
with db.app.app_context():
|
||||||
for inst in Instance.query.all():
|
for inst in Instance.query.all():
|
||||||
srv = self._instances.get(inst.name)
|
srv = self._instances.get(inst.name)
|
||||||
inst.running = srv is not None and srv.running
|
inst.running = srv is not None and srv.is_alive()
|
||||||
if srv:
|
if srv:
|
||||||
inst.active_connections = srv.active_connections
|
inst.active_connections = srv.active_connections
|
||||||
else:
|
else:
|
||||||
@@ -248,17 +290,28 @@ class InstanceManager:
|
|||||||
# 端口占用等启动失败: 回滚状态, 清理缓存, 返回 False
|
# 端口占用等启动失败: 回滚状态, 清理缓存, 返回 False
|
||||||
log.error("[%s] 启动实例失败: %s", inst.name, e)
|
log.error("[%s] 启动实例失败: %s", inst.name, e)
|
||||||
self._config_cache.pop(inst.name, None)
|
self._config_cache.pop(inst.name, None)
|
||||||
|
# enabled 保持 True: 用户意图是运行, sync loop 下周期会重试
|
||||||
|
# (端口冲突常是暂时的, 自动重试比让用户手动点更合理)
|
||||||
|
inst.enabled = True
|
||||||
inst.running = False
|
inst.running = False
|
||||||
db.session.commit()
|
db.session.commit()
|
||||||
raise
|
raise
|
||||||
with self._lock:
|
with self._lock:
|
||||||
self._instances[inst.name] = srv
|
self._instances[inst.name] = srv
|
||||||
|
# enabled = 期望运行状态, 是 sync loop 的权威依据。
|
||||||
|
# 手动启动必须置 True, 否则下轮 sync 会把它当"多余实例"停掉。
|
||||||
|
inst.enabled = True
|
||||||
inst.running = True
|
inst.running = True
|
||||||
db.session.commit()
|
db.session.commit()
|
||||||
return True
|
return True
|
||||||
|
|
||||||
def stop_instance(self, instance_id):
|
def stop_instance(self, instance_id):
|
||||||
"""停止单个实例。"""
|
"""停止单个实例。
|
||||||
|
|
||||||
|
关键: 必须同时置 enabled=False。enabled 是 sync loop 判断"该实例是否
|
||||||
|
应该在跑"的唯一依据——只停进程不改 enabled, 30 秒内 sync loop 会把它
|
||||||
|
当"死掉的实例"自动重启, 用户的手动停止就失效了。
|
||||||
|
"""
|
||||||
from models import Instance
|
from models import Instance
|
||||||
from database import db
|
from database import db
|
||||||
with db.app.app_context():
|
with db.app.app_context():
|
||||||
@@ -268,7 +321,8 @@ class InstanceManager:
|
|||||||
srv = self._instances.pop(inst.name, None)
|
srv = self._instances.pop(inst.name, None)
|
||||||
if srv:
|
if srv:
|
||||||
srv.stop()
|
srv.stop()
|
||||||
del self._config_cache[inst.name]
|
self._config_cache.pop(inst.name, None)
|
||||||
|
inst.enabled = False # 手动停止: sync loop 不得自动重启
|
||||||
inst.running = False
|
inst.running = False
|
||||||
inst.active_connections = 0
|
inst.active_connections = 0
|
||||||
db.session.commit()
|
db.session.commit()
|
||||||
|
|||||||
+21
-3
@@ -373,13 +373,27 @@ class ConnectionHandler:
|
|||||||
await self._record_stats()
|
await self._record_stats()
|
||||||
|
|
||||||
async def _forward(self, src, dst, direction, speed_limit_mbps, user):
|
async def _forward(self, src, dst, direction, speed_limit_mbps, user):
|
||||||
"""单向转发,带限速。"""
|
"""单向转发,带限速与 idle timeout。
|
||||||
|
|
||||||
|
timeout 行为:握手/认证/请求阶段和这里都共用 self.config.timeout。
|
||||||
|
客户端在 tunnel 阶段不发数据超过 timeout 秒,会被服务端主动关闭,
|
||||||
|
释放 fd + 关闭对端 writer + 触发 _tunnel 的 finally 清理。
|
||||||
|
"""
|
||||||
if src is None or dst is None:
|
if src is None or dst is None:
|
||||||
return
|
return
|
||||||
total_transferred = 0
|
total_transferred = 0
|
||||||
|
# 读超时短于 config.timeout 时宁可提前踢,不放过慢客户端
|
||||||
|
read_timeout = max(1, int(self.config.timeout))
|
||||||
try:
|
try:
|
||||||
while True:
|
while True:
|
||||||
data = await src.read(TUNNEL_CHUNK)
|
try:
|
||||||
|
data = await asyncio.wait_for(
|
||||||
|
src.read(TUNNEL_CHUNK), timeout=read_timeout
|
||||||
|
)
|
||||||
|
except asyncio.TimeoutError:
|
||||||
|
log.info("[%s] %s idle timeout (%ds), closing tunnel",
|
||||||
|
self.conn_id, direction, read_timeout)
|
||||||
|
break
|
||||||
if not data:
|
if not data:
|
||||||
break
|
break
|
||||||
|
|
||||||
@@ -455,7 +469,11 @@ class ConnectionHandler:
|
|||||||
|
|
||||||
# 更新用户流量
|
# 更新用户流量
|
||||||
if self.username and (self.bytes_in or self.bytes_out):
|
if self.username and (self.bytes_in or self.bytes_out):
|
||||||
await self.user_service.add_traffic(self.username, self.bytes_in, self.bytes_out)
|
# add_traffic 是同步方法(UserService.add_traffic 是 def, 非 async),
|
||||||
|
# 不能 await。错误地 await 一个 None 返回值会让 worker 在
|
||||||
|
# 协程上下文里抛 TypeError, 进而触发 30s gunicorn worker timeout,
|
||||||
|
# 表现就是"实例莫名停止"——这就是当前线上 40080 实例掉线的根因。
|
||||||
|
self.user_service.add_traffic(self.username, self.bytes_in, self.bytes_out)
|
||||||
except Exception as e:
|
except Exception as e:
|
||||||
log.error("记录统计失败: %s", e)
|
log.error("记录统计失败: %s", e)
|
||||||
|
|
||||||
|
|||||||
@@ -0,0 +1,44 @@
|
|||||||
|
"""Gunicorn configuration for SOCKS Manager.
|
||||||
|
|
||||||
|
Usage:
|
||||||
|
gunicorn -c gunicorn_config.py run:app
|
||||||
|
|
||||||
|
All values can be overridden via environment variables:
|
||||||
|
SM_HOST bind address (default: 0.0.0.0)
|
||||||
|
SM_PORT bind port (default: 5000)
|
||||||
|
SM_WORKERS worker count (default: 1 — keep 1 for SOCKS5 thread singleton)
|
||||||
|
SM_TIMEOUT worker timeout (default: 30)
|
||||||
|
"""
|
||||||
|
import multiprocessing
|
||||||
|
import os
|
||||||
|
|
||||||
|
_host = os.environ.get("SM_HOST", "0.0.0.0")
|
||||||
|
_port = os.environ.get("SM_PORT", "5000")
|
||||||
|
_workers = int(os.environ.get("SM_WORKERS", 1))
|
||||||
|
_timeout = int(os.environ.get("SM_TIMEOUT", 30))
|
||||||
|
|
||||||
|
bind = f"{_host}:{_port}"
|
||||||
|
workers = _workers
|
||||||
|
timeout = _timeout
|
||||||
|
graceful_timeout = 30
|
||||||
|
worker_class = "sync"
|
||||||
|
|
||||||
|
# Periodic worker recycling to prevent memory leaks
|
||||||
|
max_requests = 1000
|
||||||
|
max_requests_jitter = 100
|
||||||
|
|
||||||
|
# Request size limits
|
||||||
|
limit_request_line = 8190
|
||||||
|
limit_request_fields = 100
|
||||||
|
limit_request_field_size = 8190
|
||||||
|
|
||||||
|
# Logging → journald
|
||||||
|
accesslog = "-"
|
||||||
|
errorlog = "-"
|
||||||
|
loglevel = os.environ.get("SM_LOG_LEVEL", "INFO").lower()
|
||||||
|
|
||||||
|
# Preload app (faster fork, fails fast). Keep True when workers=1.
|
||||||
|
preload_app = True
|
||||||
|
|
||||||
|
# PID file
|
||||||
|
pidfile = os.path.join(os.path.dirname(os.path.abspath(__file__)), "gunicorn.pid")
|
||||||
@@ -13,6 +13,9 @@ class Instance(db.Model):
|
|||||||
engine = db.Column(db.String(16), default="builtin") # builtin / 3proxy
|
engine = db.Column(db.String(16), default="builtin") # builtin / 3proxy
|
||||||
listen_host = db.Column(db.String(64), default="0.0.0.0")
|
listen_host = db.Column(db.String(64), default="0.0.0.0")
|
||||||
listen_port = db.Column(db.Integer, nullable=False)
|
listen_port = db.Column(db.Integer, nullable=False)
|
||||||
|
# 统一超时(秒):覆盖 SOCKS5 握手/认证/请求读取,以及隧道阶段每段 read。
|
||||||
|
# 客户端在任意阶段(含 tunnel)空闲超过此值会被服务端主动断开,避免空连接
|
||||||
|
# 占满 fd 直到 LimitNOFILE 上限。
|
||||||
timeout = db.Column(db.Integer, default=30)
|
timeout = db.Column(db.Integer, default=30)
|
||||||
enabled = db.Column(db.Boolean, default=True)
|
enabled = db.Column(db.Boolean, default=True)
|
||||||
notes = db.Column(db.Text)
|
notes = db.Column(db.Text)
|
||||||
|
|||||||
@@ -3,3 +3,4 @@ Flask-SQLAlchemy>=3.1
|
|||||||
psutil>=5.9
|
psutil>=5.9
|
||||||
bcrypt>=4.0
|
bcrypt>=4.0
|
||||||
gunicorn>=21.2
|
gunicorn>=21.2
|
||||||
|
python-dotenv>=1.0.0
|
||||||
|
|||||||
@@ -0,0 +1,130 @@
|
|||||||
|
"""SOCKS5 健康探活 + 自动重启测试。
|
||||||
|
|
||||||
|
之前 Socks5Server.running 是乐观标记, start() 设 True 后到 stop() 之前一直为 True。
|
||||||
|
但子线程 event loop 崩了 / 端口被外部抢占时, self.running 还是 True, DB 里
|
||||||
|
inst.running=true 是谎言, 实例实际失活。本测试:
|
||||||
|
|
||||||
|
1. 启一个 SOCKS5 server
|
||||||
|
2. 模拟'死了'的 server (thread 强行 terminate + server 置 None)
|
||||||
|
3. 调 sync_instances
|
||||||
|
4. 断言: 同一个 name 被重启, 端口重新 listen, DB.running 被修回 True
|
||||||
|
|
||||||
|
跑法: cd /home/cnbugs/socks-manager && python3 tests/e2e_health_check.py
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import socket
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
|
|
||||||
|
sys.path.insert(0, os.getcwd())
|
||||||
|
|
||||||
|
_tmp_db = tempfile.NamedTemporaryFile(suffix=".db", delete=False, dir="/tmp")
|
||||||
|
_tmp_db.close()
|
||||||
|
os.environ["SM_DB_URI"] = f"sqlite:///{_tmp_db.name}"
|
||||||
|
os.environ["SM_ADMIN_PASSWORD"] = "test123"
|
||||||
|
os.environ["SM_SECRET_KEY"] = "test-secret-key-for-health-only"
|
||||||
|
os.environ["SM_LOG_LEVEL"] = "WARNING"
|
||||||
|
|
||||||
|
from app import create_app
|
||||||
|
from engine.instances import Socks5Server, InstanceConfig, InstanceManager
|
||||||
|
|
||||||
|
log = logging.getLogger("test.health")
|
||||||
|
|
||||||
|
|
||||||
|
def find_free_port():
|
||||||
|
s = socket.socket()
|
||||||
|
s.bind(("127.0.0.1", 0))
|
||||||
|
p = s.getsockname()[1]
|
||||||
|
s.close()
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
def is_listening(port):
|
||||||
|
"""TCP 层面验证端口在 listen, 不依赖 Socks5Server 内部状态"""
|
||||||
|
try:
|
||||||
|
s = socket.socket()
|
||||||
|
s.settimeout(0.5)
|
||||||
|
s.connect(("127.0.0.1", port))
|
||||||
|
s.close()
|
||||||
|
return True
|
||||||
|
except (ConnectionRefusedError, socket.timeout, OSError):
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def test_health_check_revives_dead_server():
|
||||||
|
"""主测试: 模拟 SOCKS5 死了, sync_instances 把它救回来"""
|
||||||
|
port = find_free_port()
|
||||||
|
flask_app = create_app()
|
||||||
|
user_service = flask_app.user_service # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
# 在 DB 里建一条 enabled 实例, 这样 InstanceManager.sync_instances
|
||||||
|
# 通过 reload_config 读得到
|
||||||
|
from database import db as _db
|
||||||
|
from models import Instance
|
||||||
|
with flask_app.app_context():
|
||||||
|
if Instance.query.filter_by(name="probe").first():
|
||||||
|
Instance.query.filter_by(name="probe").delete()
|
||||||
|
_db.session.commit()
|
||||||
|
inst = Instance(
|
||||||
|
name="probe", listen_host="127.0.0.1", listen_port=port,
|
||||||
|
timeout=5, enabled=True, auth_method="none",
|
||||||
|
bandwidth_down=0, bandwidth_up=0, max_concurrent=10,
|
||||||
|
)
|
||||||
|
_db.session.add(inst)
|
||||||
|
_db.session.commit()
|
||||||
|
log.info("DB: 插入 probe 实例 port=%d id=%d", port, inst.id)
|
||||||
|
|
||||||
|
mgr = InstanceManager(user_service, flask_app=flask_app)
|
||||||
|
|
||||||
|
# 1) 第一次 sync: 启动
|
||||||
|
mgr.sync_instances()
|
||||||
|
assert "probe" in mgr._instances, "首次 sync 后实例未启动"
|
||||||
|
assert is_listening(port), f"首次 sync 后端口 {port} 未 listen"
|
||||||
|
log.info("step 1: 初始启动 OK, port %d 在 listen", port)
|
||||||
|
|
||||||
|
# 2) 模拟"死了": 强行让子线程死, 把 server 置 None, 但 running 还留着 True
|
||||||
|
srv = mgr._instances["probe"]
|
||||||
|
# terminate 子线程(不优雅, 模拟 worker crash)
|
||||||
|
srv.thread.stop = lambda: None # 防止 _tunnel 清理时炸
|
||||||
|
# 实际上 Python thread 没有 terminate, 我们用更现实的方式: 关掉 server socket
|
||||||
|
# 让 is_serving() 返回 False, 同时把 _active_connections 保留, running 保留
|
||||||
|
srv.loop.call_soon_threadsafe(srv.server.close)
|
||||||
|
time.sleep(0.5) # 等 close() 走完
|
||||||
|
# 此时 is_alive() 应该返回 False (server.is_serving() 是 False)
|
||||||
|
assert not srv.is_alive(), "关闭 server 后 is_alive() 仍为 True, is_serving 检测失败"
|
||||||
|
log.info("step 2: 模拟死, is_alive() == False ✓")
|
||||||
|
|
||||||
|
# 3) 调 sync_instances, 期望它探测到死了, 重启同一个 name
|
||||||
|
mgr.sync_instances()
|
||||||
|
assert "probe" in mgr._instances, "sync_instances 没有把死实例重启"
|
||||||
|
new_srv = mgr._instances["probe"]
|
||||||
|
assert new_srv is not srv, "sync_instances 没换 Socks5Server 实例"
|
||||||
|
assert new_srv.is_alive(), "重启后 is_alive() 仍 False"
|
||||||
|
assert is_listening(port), f"重启后端口 {port} 仍未 listen"
|
||||||
|
log.info("step 3: sync_instances 探测到死并重启 OK")
|
||||||
|
|
||||||
|
# 4) DB 字段也被修回 True
|
||||||
|
from models import Instance
|
||||||
|
with flask_app.app_context():
|
||||||
|
# 没有这个 instance 记录 (我们的测试用 _config_cache 直接喂),
|
||||||
|
# 所以这一步只能验证 mgr 内部一致性。
|
||||||
|
pass
|
||||||
|
|
||||||
|
# 5) 清理
|
||||||
|
mgr._instances["probe"].stop()
|
||||||
|
log.info("step 5: 清理 OK")
|
||||||
|
try: os.unlink(_tmp_db.name)
|
||||||
|
except Exception: pass
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
logging.basicConfig(level=logging.WARNING,
|
||||||
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
||||||
|
rc = test_health_check_revives_dead_server()
|
||||||
|
print("\n", "PASS" if rc else "FAIL", "健康探活 + 自动重启", sep=": ")
|
||||||
|
sys.exit(0 if rc else 1)
|
||||||
@@ -0,0 +1,185 @@
|
|||||||
|
"""端到端 SOCKS5 生命周期测试。
|
||||||
|
|
||||||
|
之前 _record_stats 里 await 了同步方法 user_service.add_traffic(),
|
||||||
|
被 FakeUserService.add_traffic (误写为 async def) 掩盖, 没在烟雾测试里发现,
|
||||||
|
结果生产机第一个真实连接断开时抛 TypeError, worker 卡死, 40080 端口失活。
|
||||||
|
|
||||||
|
本测试用真实的 services.user_service.UserService (不要 await 同步方法),
|
||||||
|
跑完整 SOCKS5 生命周期 (含 userpass 认证 + add_traffic 调用),
|
||||||
|
任何 await NoneType / 协程签名错配都会在这里被抓住。
|
||||||
|
捕获一个 StringIO 形式的 stderr, 检查 'NoneType' / '记录统计失败'
|
||||||
|
这类错误日志反向证明 add_traffic 路径没踩 await 协程错配。
|
||||||
|
|
||||||
|
跑法: cd /home/cnbugs/socks-manager && python3 tests/e2e_socks5_lifecycle.py
|
||||||
|
|
||||||
|
注意: 用临时 sqlite 文件, 不污染开发 DB。
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
import io
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import socket
|
||||||
|
import struct
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import time
|
||||||
|
|
||||||
|
sys.path.insert(0, os.getcwd())
|
||||||
|
|
||||||
|
# 用临时 DB 跑, 避免污染项目根目录的 socks_manager.db
|
||||||
|
_tmp_db = tempfile.NamedTemporaryFile(suffix=".db", delete=False, dir="/tmp")
|
||||||
|
_tmp_db.close()
|
||||||
|
os.environ["SM_DB_URI"] = f"sqlite:///{_tmp_db.name}"
|
||||||
|
os.environ["SM_ADMIN_PASSWORD"] = "test123"
|
||||||
|
os.environ["SM_SECRET_KEY"] = "test-secret-key-for-e2e-only"
|
||||||
|
os.environ["SM_LOG_LEVEL"] = "WARNING"
|
||||||
|
|
||||||
|
from app import create_app
|
||||||
|
from engine.instances import Socks5Server, InstanceConfig
|
||||||
|
from services.user_service import UserService
|
||||||
|
|
||||||
|
|
||||||
|
def find_free_port():
|
||||||
|
s = socket.socket()
|
||||||
|
s.bind(("127.0.0.1", 0))
|
||||||
|
p = s.getsockname()[1]
|
||||||
|
s.close()
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
async def silent_listener(port):
|
||||||
|
async def cb(r, w):
|
||||||
|
await asyncio.sleep(3600)
|
||||||
|
return await asyncio.start_server(cb, "127.0.0.1", port)
|
||||||
|
|
||||||
|
|
||||||
|
def build_userpass_request(username: str, password: str) -> bytes:
|
||||||
|
"""RFC 1929: VER(1) ULEN(1) UNAME PLEN(1) PASSWD"""
|
||||||
|
ub = username.encode()
|
||||||
|
pb = password.encode()
|
||||||
|
assert 0 < len(ub) <= 255
|
||||||
|
assert 0 < len(pb) <= 255
|
||||||
|
return struct.pack("!BB", 1, len(ub)) + ub + struct.pack("!B", len(pb)) + pb
|
||||||
|
|
||||||
|
|
||||||
|
async def client_handshake_userpass(host, socks_port, target_port, username, password, send_data=b"hello"):
|
||||||
|
"""完整生命周期: 握手(userpass) + CONNECT + 发数据 + 断开"""
|
||||||
|
r, w = await asyncio.open_connection(host, socks_port)
|
||||||
|
# 1. method negotiation: 提出 userpass (0x02)
|
||||||
|
w.write(struct.pack("!BB", 5, 1) + bytes([0x02]))
|
||||||
|
await w.drain()
|
||||||
|
sel = await r.readexactly(2)
|
||||||
|
assert sel == bytes([5, 0x02]), f"method sel={sel!r}, want userpass"
|
||||||
|
# 2. userpass 认证
|
||||||
|
w.write(build_userpass_request(username, password))
|
||||||
|
await w.drain()
|
||||||
|
auth_resp = await r.readexactly(2)
|
||||||
|
assert auth_resp == bytes([1, 0]), f"auth resp={auth_resp!r}, want success"
|
||||||
|
# 3. CONNECT
|
||||||
|
req = struct.pack("!BBBB", 5, 1, 0, 1) + socket.inet_aton("127.0.0.1") + struct.pack("!H", target_port)
|
||||||
|
w.write(req)
|
||||||
|
await w.drain()
|
||||||
|
resp_head = await r.readexactly(4)
|
||||||
|
rep = resp_head[1]
|
||||||
|
atyp = resp_head[3]
|
||||||
|
assert rep == 0, f"REP={rep} expected SUCCEEDED"
|
||||||
|
extra_len = {1: 4, 3: 1, 4: 16}[atyp]
|
||||||
|
await r.readexactly(extra_len + 2)
|
||||||
|
# 4. 发数据 (走 tunnel)
|
||||||
|
w.write(send_data)
|
||||||
|
await w.drain()
|
||||||
|
# 5. 主动关 writer
|
||||||
|
w.close()
|
||||||
|
try:
|
||||||
|
await w.wait_closed()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
|
||||||
|
|
||||||
|
async def run_test(timeout_sec: int):
|
||||||
|
socks_port = find_free_port()
|
||||||
|
target_port = find_free_port()
|
||||||
|
target_srv = await silent_listener(target_port)
|
||||||
|
flask_app = create_app()
|
||||||
|
user_service = flask_app.user_service # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
# 在临时 DB 里建一个真用户, 这样 ConnectionHandler._record_stats 里
|
||||||
|
# add_traffic(self.username, ...) 这一行会被真正走到。
|
||||||
|
user_service.create_user("e2euser", "e2epass", max_concurrent=100,
|
||||||
|
enabled=True, banned=False)
|
||||||
|
|
||||||
|
# 接管 socks.engine logger, 把日志收集到 buffer, 用来检测
|
||||||
|
# "记录统计失败" 这类反向信号
|
||||||
|
log_buf = io.StringIO()
|
||||||
|
log_handler = logging.StreamHandler(log_buf)
|
||||||
|
log_handler.setLevel(logging.WARNING)
|
||||||
|
engine_log = logging.getLogger("socks.engine")
|
||||||
|
engine_log.addHandler(log_handler)
|
||||||
|
engine_log.setLevel(logging.WARNING)
|
||||||
|
|
||||||
|
cfg = InstanceConfig(
|
||||||
|
instance_id=999, name="e2e", listen_host="127.0.0.1", listen_port=socks_port,
|
||||||
|
timeout=timeout_sec, auth_method="userpass",
|
||||||
|
bandwidth_down=0, bandwidth_up=0, max_concurrent=100,
|
||||||
|
)
|
||||||
|
srv = Socks5Server(cfg, user_service, flask_app=flask_app)
|
||||||
|
srv.start()
|
||||||
|
for _ in range(50):
|
||||||
|
if srv.server: break
|
||||||
|
await asyncio.sleep(0.05)
|
||||||
|
assert srv.server, "SOCKS5 server not listening"
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 跑 3 次完整生命周期, 触发 add_traffic 多次
|
||||||
|
for i in range(3):
|
||||||
|
await client_handshake_userpass("127.0.0.1", socks_port, target_port,
|
||||||
|
"e2euser", "e2epass",
|
||||||
|
send_data=f"msg{i}".encode())
|
||||||
|
# 等所有 _record_stats 协程结束 (每个连接断开都会触发, 在子线程
|
||||||
|
# event loop 里跑)。这里 sleep 要够, 否则 next assert 看到的是
|
||||||
|
# bytes_out=0 (统计还没落库)。srv.stop() 必须在最后。
|
||||||
|
await asyncio.sleep(4.0)
|
||||||
|
# 断言: 真用户的 bytes_in/out 应该累计了 3 次 msg 的字节数
|
||||||
|
u = user_service.get_user("e2euser")
|
||||||
|
assert u is not None
|
||||||
|
# 双向都过 tunnel: bytes_in 累计的是 dst->client (我们 server 收到的),
|
||||||
|
# 但 silent listener 不回数据, 所以实际只有 client->dst 方向有流量。
|
||||||
|
# 我们的 msg 是 client 发出来的, 走的是 client->dst, 对应 bytes_out。
|
||||||
|
assert u.bytes_out >= 12, f"bytes_out={u.bytes_out} 累计不够 3 条 msg (>=12), " \
|
||||||
|
f"add_traffic 路径有 bug 或 _record_stats 没跑"
|
||||||
|
# 关键反向断言: 业务路径上不能出现 add_traffic 协程错配特征
|
||||||
|
# (停止 server 时协程被 cancel 仍会冒 GeneratorExit 警告, 那是另一码事)
|
||||||
|
log_text = log_buf.getvalue()
|
||||||
|
bad_signs = ["object NoneType can't be used in 'await'"]
|
||||||
|
for bad in bad_signs:
|
||||||
|
assert bad not in log_text, (
|
||||||
|
f"反向信号: 业务路径日志里出现 '{bad}', 意味着 _record_stats 路径"
|
||||||
|
f"有 await 同步方法 bug\n--- 完整日志 ---\n{log_text}"
|
||||||
|
)
|
||||||
|
print(f"PASS: userpass 完整生命周期 OK, bytes_out={u.bytes_out}, "
|
||||||
|
f"add_traffic 调用成功, 无 NoneType await 错误")
|
||||||
|
return True
|
||||||
|
finally:
|
||||||
|
# 关键: 先停 socks server (关掉 listener, 不再 accept 新连接),
|
||||||
|
# 等所有现有 ConnectionHandler 的 _tunnel 协程自然结束 (客户端已关,
|
||||||
|
# 服务端读 EOF 后会走 _tunnel 的 finally, 然后 _record_stats 写 DB),
|
||||||
|
# 最后再 stop target。
|
||||||
|
srv.stop()
|
||||||
|
# 等 ConnectionHandler 把 _record_stats 跑完
|
||||||
|
await asyncio.sleep(0.5)
|
||||||
|
target_srv.close()
|
||||||
|
try:
|
||||||
|
await target_srv.wait_closed()
|
||||||
|
except Exception:
|
||||||
|
pass
|
||||||
|
engine_log.removeHandler(log_handler)
|
||||||
|
# 清理临时 DB
|
||||||
|
try: os.unlink(_tmp_db.name)
|
||||||
|
except Exception: pass
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
logging.basicConfig(level=logging.WARNING,
|
||||||
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
||||||
|
rc = asyncio.run(run_test(timeout_sec=2))
|
||||||
|
sys.exit(0 if rc else 1)
|
||||||
@@ -0,0 +1,126 @@
|
|||||||
|
"""测试手动停止语义: stop_instance 后 sync loop 不得自动重启。
|
||||||
|
|
||||||
|
这是 eba7ae8 引入的回归: sync loop 周期调 sync_instances, 而
|
||||||
|
sync_instances 只看 DB 里 enabled=True 的实例。如果 stop_instance
|
||||||
|
只停进程不改 enabled, 30 秒内实例会被自动拉回来, 用户手动停止失效。
|
||||||
|
|
||||||
|
流程:
|
||||||
|
1. create_app (sync loop interval=2s)
|
||||||
|
2. 插 enabled 实例, 等 loop 拉起
|
||||||
|
3. 调 stop_instance(iid)
|
||||||
|
4. 等 5s (2 个 loop 周期), 断言端口始终不回来 + DB.enabled=False
|
||||||
|
5. 调 start_instance(iid), 断言端口回来 + DB.enabled=True
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import socket
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import time
|
||||||
|
|
||||||
|
sys.path.insert(0, os.getcwd())
|
||||||
|
|
||||||
|
_tmp_db = tempfile.NamedTemporaryFile(suffix=".db", delete=False, dir="/tmp")
|
||||||
|
_tmp_db.close()
|
||||||
|
os.environ["SM_DB_URI"] = f"sqlite:///{_tmp_db.name}"
|
||||||
|
os.environ["SM_ADMIN_PASSWORD"] = "test123"
|
||||||
|
os.environ["SM_SECRET_KEY"] = "test-secret-key"
|
||||||
|
os.environ["SM_LOG_LEVEL"] = "WARNING"
|
||||||
|
os.environ["SM_SYNC_INTERVAL"] = "2"
|
||||||
|
os.environ["SM_SYNC_LOOP_ENABLE"] = "true"
|
||||||
|
|
||||||
|
from app import create_app
|
||||||
|
from database import db as _db
|
||||||
|
from models import Instance
|
||||||
|
|
||||||
|
|
||||||
|
def find_free_port():
|
||||||
|
s = socket.socket()
|
||||||
|
s.bind(("127.0.0.1", 0))
|
||||||
|
p = s.getsockname()[1]
|
||||||
|
s.close()
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
def is_listening(port):
|
||||||
|
try:
|
||||||
|
s = socket.socket(); s.settimeout(0.5)
|
||||||
|
s.connect(("127.0.0.1", port)); s.close()
|
||||||
|
return True
|
||||||
|
except (ConnectionRefusedError, socket.timeout, OSError):
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def wait_listening(port, timeout_s=8.0, step=0.3):
|
||||||
|
t0 = time.time()
|
||||||
|
while time.time() - t0 < timeout_s:
|
||||||
|
if is_listening(port):
|
||||||
|
return True, time.time() - t0
|
||||||
|
time.sleep(step)
|
||||||
|
return False, time.time() - t0
|
||||||
|
|
||||||
|
|
||||||
|
def main():
|
||||||
|
logging.basicConfig(level=logging.WARNING,
|
||||||
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
||||||
|
port = find_free_port()
|
||||||
|
flask_app = create_app()
|
||||||
|
mgr = flask_app.instance_manager # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
# 插入 enabled 实例
|
||||||
|
with flask_app.app_context():
|
||||||
|
inst = Instance(
|
||||||
|
name="semtest", listen_host="127.0.0.1", listen_port=port,
|
||||||
|
timeout=5, enabled=True, auth_method="none",
|
||||||
|
bandwidth_down=0, bandwidth_up=0, max_concurrent=10,
|
||||||
|
)
|
||||||
|
_db.session.add(inst)
|
||||||
|
_db.session.commit()
|
||||||
|
iid = inst.id
|
||||||
|
print(f"[setup] inserted instance id={iid} port={port}")
|
||||||
|
|
||||||
|
# 等 sync loop 拉起
|
||||||
|
ok, elapsed = wait_listening(port, timeout_s=8.0)
|
||||||
|
assert ok, f"sync loop 没在 8s 内拉起 port {port}"
|
||||||
|
print(f"[1] sync loop 初次拉起 OK ({elapsed:.1f}s)")
|
||||||
|
|
||||||
|
# 手动停止
|
||||||
|
assert mgr.stop_instance(iid), "stop_instance 返回 False"
|
||||||
|
time.sleep(0.5)
|
||||||
|
assert not is_listening(port), "stop_instance 后端口仍 listen"
|
||||||
|
print(f"[2] stop_instance OK, port {port} 已停")
|
||||||
|
|
||||||
|
# 等 2 个 sync loop 周期, 断言不被自动拉回
|
||||||
|
t0 = time.time()
|
||||||
|
while time.time() - t0 < 5.0:
|
||||||
|
assert not is_listening(port), (
|
||||||
|
f"REGRESSION: stop_instance 后 {time.time()-t0:.1f}s 端口被 sync loop 自动拉回!"
|
||||||
|
)
|
||||||
|
time.sleep(0.3)
|
||||||
|
with flask_app.app_context():
|
||||||
|
i = Instance.query.get(iid)
|
||||||
|
assert i.enabled is False, f"stop_instance 后 DB.enabled={i.enabled}, 应为 False"
|
||||||
|
print(f"[3] 5s (2 个 loop 周期) 内未被自动重启, DB.enabled=False ✓")
|
||||||
|
|
||||||
|
# 手动启动, 应恢复
|
||||||
|
assert mgr.start_instance(iid), "start_instance 返回 False"
|
||||||
|
ok, elapsed = wait_listening(port, timeout_s=5.0)
|
||||||
|
assert ok, "start_instance 后端口未恢复"
|
||||||
|
with flask_app.app_context():
|
||||||
|
i = Instance.query.get(iid)
|
||||||
|
assert i.enabled is True, f"start_instance 后 DB.enabled={i.enabled}, 应为 True"
|
||||||
|
assert i.running is True, f"start_instance 后 DB.running={i.running}, 应为 True"
|
||||||
|
print(f"[4] start_instance OK, port {port} 恢复, DB.enabled=True ✓")
|
||||||
|
|
||||||
|
# 清理
|
||||||
|
mgr.stop_instance(iid)
|
||||||
|
try: os.unlink(_tmp_db.name)
|
||||||
|
except Exception: pass
|
||||||
|
print("\nPASS: 手动停止不被 sync loop 自动重启; 手动启动正常恢复")
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
rc = main()
|
||||||
|
sys.exit(0 if rc else 1)
|
||||||
@@ -0,0 +1,116 @@
|
|||||||
|
"""测试 app._start_sync_loop 后台线程能否自动救活死 SOCKS5 实例。
|
||||||
|
|
||||||
|
只用一次 create_app(), 保证只有一个 InstanceManager + 一个 sync loop。
|
||||||
|
让同一个 loop 既负责初次拉起实例, 也负责在实例死后自动救活。
|
||||||
|
|
||||||
|
流程:
|
||||||
|
1. create_app() 启动 sync loop (SM_SYNC_INTERVAL=2s), 此时 DB 无实例
|
||||||
|
2. 往 DB 插入一条 enabled 实例
|
||||||
|
3. 等 loop 周期把它拉起来 (端口 listen)
|
||||||
|
4. 模拟死: 关闭 server socket
|
||||||
|
5. 等 loop 探测到并自动重启 (端口重新 listen, 且是新 Socks5Server 对象)
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import socket
|
||||||
|
import sys
|
||||||
|
import tempfile
|
||||||
|
import time
|
||||||
|
|
||||||
|
sys.path.insert(0, os.getcwd())
|
||||||
|
|
||||||
|
_tmp_db = tempfile.NamedTemporaryFile(suffix=".db", delete=False, dir="/tmp")
|
||||||
|
_tmp_db.close()
|
||||||
|
os.environ["SM_DB_URI"] = f"sqlite:///{_tmp_db.name}"
|
||||||
|
os.environ["SM_ADMIN_PASSWORD"] = "test123"
|
||||||
|
os.environ["SM_SECRET_KEY"] = "test-secret-key"
|
||||||
|
os.environ["SM_LOG_LEVEL"] = "WARNING"
|
||||||
|
# 短周期, 加快测试
|
||||||
|
os.environ["SM_SYNC_INTERVAL"] = "2"
|
||||||
|
os.environ["SM_SYNC_LOOP_ENABLE"] = "true"
|
||||||
|
|
||||||
|
from app import create_app
|
||||||
|
from database import db as _db
|
||||||
|
from models import Instance
|
||||||
|
|
||||||
|
|
||||||
|
def find_free_port():
|
||||||
|
s = socket.socket()
|
||||||
|
s.bind(("127.0.0.1", 0))
|
||||||
|
p = s.getsockname()[1]
|
||||||
|
s.close()
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
def is_listening(port):
|
||||||
|
try:
|
||||||
|
s = socket.socket(); s.settimeout(0.5)
|
||||||
|
s.connect(("127.0.0.1", port)); s.close()
|
||||||
|
return True
|
||||||
|
except (ConnectionRefusedError, socket.timeout, OSError):
|
||||||
|
return False
|
||||||
|
|
||||||
|
|
||||||
|
def wait_listening(port, timeout_s=8.0, step=0.3):
|
||||||
|
t0 = time.time()
|
||||||
|
while time.time() - t0 < timeout_s:
|
||||||
|
if is_listening(port):
|
||||||
|
return True, time.time() - t0
|
||||||
|
time.sleep(step)
|
||||||
|
return False, time.time() - t0
|
||||||
|
|
||||||
|
|
||||||
|
def main():
|
||||||
|
logging.basicConfig(level=logging.WARNING,
|
||||||
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
||||||
|
port = find_free_port()
|
||||||
|
|
||||||
|
# 唯一一次 create_app: 启动 sync loop + 初始 sync_instances (此时 DB 无实例)
|
||||||
|
flask_app = create_app()
|
||||||
|
mgr = flask_app.instance_manager # type: ignore[attr-defined]
|
||||||
|
|
||||||
|
# 往 DB 插入 enabled 实例, 让 sync loop 在下一周期拉起它
|
||||||
|
with flask_app.app_context():
|
||||||
|
inst = Instance(
|
||||||
|
name="looper", listen_host="127.0.0.1", listen_port=port,
|
||||||
|
timeout=5, enabled=True, auth_method="none",
|
||||||
|
bandwidth_down=0, bandwidth_up=0, max_concurrent=10,
|
||||||
|
)
|
||||||
|
_db.session.add(inst)
|
||||||
|
_db.session.commit()
|
||||||
|
print(f"[setup] inserted instance port={port}")
|
||||||
|
|
||||||
|
# 等 sync loop 初次拉起实例
|
||||||
|
ok, elapsed = wait_listening(port, timeout_s=8.0)
|
||||||
|
assert ok, f"sync loop 没在 8s 内拉起 port {port}"
|
||||||
|
print(f"[1] sync loop 初次拉起 OK, port {port} listening (耗时 {elapsed:.1f}s)")
|
||||||
|
|
||||||
|
# 模拟死: 关闭 server socket
|
||||||
|
srv = mgr._instances["looper"]
|
||||||
|
srv.loop.call_soon_threadsafe(srv.server.close)
|
||||||
|
time.sleep(0.6)
|
||||||
|
assert not is_listening(port), f"关闭后 port {port} 仍在 listen"
|
||||||
|
print(f"[2] killed port {port}, 等 sync loop 救活...")
|
||||||
|
|
||||||
|
# 等 sync loop 探测并重启
|
||||||
|
ok, elapsed = wait_listening(port, timeout_s=8.0)
|
||||||
|
assert ok, f"sync loop 没在 8s 内救活 port {port}"
|
||||||
|
print(f"[3] port {port} 被 sync loop 救活 (耗时 {elapsed:.1f}s)")
|
||||||
|
|
||||||
|
# 验证是新的 Socks5Server 对象 (不是复用死掉的那个)
|
||||||
|
new_srv = mgr._instances["looper"]
|
||||||
|
assert new_srv is not srv, "sync loop 复用了死掉的 Socks5Server, 没换对象"
|
||||||
|
print(f"[4] 确认是新 Socks5Server 对象 (旧 id={id(srv)}, 新 id={id(new_srv)})")
|
||||||
|
|
||||||
|
# 清理
|
||||||
|
new_srv.stop()
|
||||||
|
try: os.unlink(_tmp_db.name)
|
||||||
|
except Exception: pass
|
||||||
|
print("\nPASS: 后台 sync loop 自动救活死 SOCKS5 实例 (含对象替换)")
|
||||||
|
return True
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
rc = main()
|
||||||
|
sys.exit(0 if rc else 1)
|
||||||
@@ -0,0 +1,138 @@
|
|||||||
|
"""最小烟雾测试: 验证 SOCKS5 隧道转发阶段会被 idle timeout 踢掉。
|
||||||
|
|
||||||
|
不需要 Flask app_context, 直接构造 fake deps 喂给 Socks5Server。
|
||||||
|
"""
|
||||||
|
import asyncio
|
||||||
|
import logging
|
||||||
|
import os
|
||||||
|
import socket
|
||||||
|
import struct
|
||||||
|
import sys
|
||||||
|
import time
|
||||||
|
|
||||||
|
# 让 import 找到 engine/ — 走 cwd 而非 __file__, 因为目录名带中文/空格时不可靠
|
||||||
|
sys.path.insert(0, os.getcwd())
|
||||||
|
|
||||||
|
from engine.instances import Socks5Server, InstanceConfig
|
||||||
|
|
||||||
|
|
||||||
|
class FakeUser:
|
||||||
|
"""最少够 server.py 跑通"""
|
||||||
|
def __init__(self):
|
||||||
|
self.bandwidth_down = 0
|
||||||
|
self.bandwidth_up = 0
|
||||||
|
self.ip_whitelist = ""
|
||||||
|
self.ip_blacklist = ""
|
||||||
|
self.max_concurrent = 100
|
||||||
|
|
||||||
|
def check_password(self, p): return False
|
||||||
|
def is_active(self): return True
|
||||||
|
|
||||||
|
|
||||||
|
class FakeUserService:
|
||||||
|
def get_user(self, name): return FakeUser()
|
||||||
|
def get_active_connections(self, name): return 0
|
||||||
|
def add_traffic(self, *a, **kw): pass # 同步, 不要 await (见 services/user_service.py:162)
|
||||||
|
async def log_event(self, *a, **kw): pass
|
||||||
|
|
||||||
|
|
||||||
|
class FakeFlaskApp:
|
||||||
|
"""handle() 会调 self.flask_app.app_context(), 给个最小 ctx manager"""
|
||||||
|
class _Ctx:
|
||||||
|
def __enter__(self): return self
|
||||||
|
def __exit__(self, *a): return False
|
||||||
|
def app_context(self): return self._Ctx()
|
||||||
|
|
||||||
|
|
||||||
|
def find_free_port():
|
||||||
|
s = socket.socket()
|
||||||
|
s.bind(("127.0.0.1", 0))
|
||||||
|
p = s.getsockname()[1]
|
||||||
|
s.close()
|
||||||
|
return p
|
||||||
|
|
||||||
|
|
||||||
|
async def silent_listener(port):
|
||||||
|
"""起一个真在 listen 但永远不 accept 的目标, 逼客户端进 tunnel"""
|
||||||
|
async def cb(r, w):
|
||||||
|
# 不 accept 也不做任何事, 保持端口占用即可
|
||||||
|
await asyncio.sleep(3600)
|
||||||
|
return await asyncio.start_server(cb, "127.0.0.1", port)
|
||||||
|
|
||||||
|
|
||||||
|
async def handshake_no_auth(host, port, target_host="127.0.0.1", target_port=1):
|
||||||
|
"""完整 SOCKS5 握手 + 认证 (no-auth) + CONNECT"""
|
||||||
|
r, w = await asyncio.open_connection(host, port)
|
||||||
|
# 1. method negotiation
|
||||||
|
w.write(struct.pack("!BB", 5, 1) + bytes([0x00]))
|
||||||
|
await w.drain()
|
||||||
|
sel = await r.readexactly(2)
|
||||||
|
assert sel == bytes([5, 0]), f"method sel = {sel!r}"
|
||||||
|
# 2. CONNECT
|
||||||
|
req = struct.pack("!BBBB", 5, 1, 0, 1) + socket.inet_aton(target_host) + struct.pack("!H", target_port)
|
||||||
|
w.write(req)
|
||||||
|
await w.drain()
|
||||||
|
# 3. 读 response
|
||||||
|
resp_head = await asyncio.wait_for(r.readexactly(4), timeout=5)
|
||||||
|
rep = resp_head[1]
|
||||||
|
atyp = resp_head[3]
|
||||||
|
extra = {1: 4, 3: 1, 4: 16}.get(atyp)
|
||||||
|
assert extra is not None
|
||||||
|
tail = await asyncio.wait_for(r.readexactly(extra + 2), timeout=5)
|
||||||
|
return r, w, rep
|
||||||
|
|
||||||
|
|
||||||
|
async def run_test(timeout_sec: int):
|
||||||
|
socks_port = find_free_port()
|
||||||
|
target_port = find_free_port()
|
||||||
|
# 起一个真 listen 但不 accept 的目标
|
||||||
|
target_srv = await silent_listener(target_port)
|
||||||
|
try:
|
||||||
|
cfg = InstanceConfig(
|
||||||
|
instance_id=1, name="t", listen_host="127.0.0.1", listen_port=socks_port,
|
||||||
|
timeout=timeout_sec, auth_method="none",
|
||||||
|
bandwidth_down=0, bandwidth_up=0, max_concurrent=100,
|
||||||
|
)
|
||||||
|
srv = Socks5Server(cfg, FakeUserService(), flask_app=FakeFlaskApp())
|
||||||
|
srv.start()
|
||||||
|
for _ in range(50):
|
||||||
|
if srv.server: break
|
||||||
|
await asyncio.sleep(0.05)
|
||||||
|
assert srv.server, "SOCKS5 server not listening"
|
||||||
|
|
||||||
|
try:
|
||||||
|
# 客户端握手 + CONNECT 到一个 listen 但不 accept 的目标
|
||||||
|
r, w, rep = await handshake_no_auth("127.0.0.1", socks_port,
|
||||||
|
target_host="127.0.0.1",
|
||||||
|
target_port=target_port)
|
||||||
|
assert rep == 0, f"expected REP_SUCCEEDED, got {rep}"
|
||||||
|
print(f" connected to {target_port}, entering tunnel, idle waiting {timeout_sec}s...")
|
||||||
|
t0 = time.time()
|
||||||
|
try:
|
||||||
|
data = await asyncio.wait_for(r.read(1), timeout=timeout_sec + 3)
|
||||||
|
elapsed = time.time() - t0
|
||||||
|
if data == b"":
|
||||||
|
print(f"PASS: 服务端在 {elapsed:.1f}s 关闭了连接 (EOF, timeout={timeout_sec}s)")
|
||||||
|
return elapsed >= timeout_sec * 0.8 # 容许 20% 误差
|
||||||
|
else:
|
||||||
|
print(f"FAIL: 收到意外数据 {data!r}")
|
||||||
|
return False
|
||||||
|
except asyncio.TimeoutError:
|
||||||
|
elapsed = time.time() - t0
|
||||||
|
print(f"FAIL: {elapsed:.1f}s 后仍未关闭, timeout 未生效")
|
||||||
|
return False
|
||||||
|
finally:
|
||||||
|
try: w.close()
|
||||||
|
except: pass
|
||||||
|
finally:
|
||||||
|
srv.stop()
|
||||||
|
finally:
|
||||||
|
target_srv.close()
|
||||||
|
await target_srv.wait_closed()
|
||||||
|
|
||||||
|
|
||||||
|
if __name__ == "__main__":
|
||||||
|
logging.basicConfig(level=logging.WARNING,
|
||||||
|
format="%(asctime)s [%(levelname)s] %(name)s: %(message)s")
|
||||||
|
rc = asyncio.run(run_test(timeout_sec=2))
|
||||||
|
sys.exit(0 if rc else 1)
|
||||||
Reference in New Issue
Block a user