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
|
||||
pip install gunicorn --break-system-packages
|
||||
|
||||
# 启动(4 worker 进程,适合 4 核服务器)
|
||||
gunicorn -w 4 -b 0.0.0.0:5000 --timeout 120 run:app
|
||||
# 最小启动(不推荐生产,worker 数量 = CPU 核心数,单核机器用 -w 1)
|
||||
gunicorn -w 1 -b 0.0.0.0:5000 run:app
|
||||
|
||||
# 推荐参数(更长超时 + 优雅关闭)
|
||||
gunicorn -w 2 -b 0.0.0.0:5000 \
|
||||
--timeout 300 \
|
||||
# 推荐参数:30s 超时 + 优雅关闭 + 周期性回收 worker 防内存泄漏
|
||||
gunicorn -w 1 -b 0.0.0.0:5000 \
|
||||
--timeout 30 \
|
||||
--graceful-timeout 30 \
|
||||
--max-requests 10000 \
|
||||
--max-requests-jitter 1000 \
|
||||
--max-requests 1000 \
|
||||
--max-requests-jitter 100 \
|
||||
--limit-request-line 8190 \
|
||||
--limit-request-fields 100 \
|
||||
run:app
|
||||
```
|
||||
|
||||
> **注意**:SOCKS5 代理实例运行在独立的 asyncio 线程中,不受 Gunicorn worker 数量影响。
|
||||
> **生产请用 systemd 托管**,见下节;裸 gunicorn 没有开机自启和崩溃重启。
|
||||
|
||||
---
|
||||
|
||||
### 2. Systemd 服务配置(推荐)
|
||||
|
||||
创建 `/etc/systemd/system/socks-manager.service`:
|
||||
服务单元文件的**唯一权威源**在仓库 `deploy/socks-manager.service`。
|
||||
所有参数、调优注释、`MemoryMax` 限制都集中在那一份,生产机只放一份拷贝。
|
||||
|
||||
```ini
|
||||
[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"
|
||||
|
||||
# 生产环境变量(按需修改)
|
||||
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
|
||||
```bash
|
||||
# 1. 在本地修改 deploy/socks-manager.service
|
||||
# 2. git commit + push
|
||||
# 3. 在生产机 (/opt/socks-manager) 拉取并执行:
|
||||
git pull
|
||||
sudo bash deploy/install.sh
|
||||
```
|
||||
|
||||
启动服务:
|
||||
`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
|
||||
systemctl daemon-reload
|
||||
systemctl enable --now 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 应用工厂。"""
|
||||
import atexit
|
||||
import os
|
||||
import logging
|
||||
import threading
|
||||
import time
|
||||
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 database import db
|
||||
from models import Instance # 触发表创建
|
||||
@@ -11,6 +24,58 @@ from engine.instances import InstanceManager
|
||||
from services.user_service import UserService
|
||||
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():
|
||||
# 确保密码已设置(首次启动自动设置默认密码)
|
||||
@@ -63,4 +128,13 @@ def create_app():
|
||||
app.register_blueprint(web_bp)
|
||||
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
|
||||
|
||||
@@ -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:
|
||||
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
|
||||
|
||||
@@ -185,7 +211,12 @@ class InstanceManager:
|
||||
self._config_cache = new_configs
|
||||
|
||||
def sync_instances(self):
|
||||
"""根据数据库配置同步实例运行状态。"""
|
||||
"""根据数据库配置同步实例运行状态。
|
||||
|
||||
健康检查: 对每个 _instances[name], 调用 is_alive() 验证线程 + server
|
||||
socket 都活着。如果死了, 从 _instances 删掉, 用同一个 config 重新启动。
|
||||
避免 40080 那种'DB 写 running=true 但端口没人 listen'的谎言状态。
|
||||
"""
|
||||
from models import Instance
|
||||
from database import db
|
||||
self.reload_config()
|
||||
@@ -193,7 +224,18 @@ class InstanceManager:
|
||||
current_names = set(self._config_cache.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:
|
||||
cfg = self._config_cache[name]
|
||||
srv = Socks5Server(cfg, self.user_service, self.flask_app)
|
||||
@@ -210,11 +252,11 @@ class InstanceManager:
|
||||
srv = self._instances.pop(name)
|
||||
srv.stop()
|
||||
|
||||
# 更新实例运行状态到数据库
|
||||
# 更新实例运行状态到数据库 (用 is_alive() 而不是乐观 running 标记)
|
||||
with db.app.app_context():
|
||||
for inst in Instance.query.all():
|
||||
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:
|
||||
inst.active_connections = srv.active_connections
|
||||
else:
|
||||
@@ -248,17 +290,28 @@ class InstanceManager:
|
||||
# 端口占用等启动失败: 回滚状态, 清理缓存, 返回 False
|
||||
log.error("[%s] 启动实例失败: %s", inst.name, e)
|
||||
self._config_cache.pop(inst.name, None)
|
||||
# enabled 保持 True: 用户意图是运行, sync loop 下周期会重试
|
||||
# (端口冲突常是暂时的, 自动重试比让用户手动点更合理)
|
||||
inst.enabled = True
|
||||
inst.running = False
|
||||
db.session.commit()
|
||||
raise
|
||||
with self._lock:
|
||||
self._instances[inst.name] = srv
|
||||
# enabled = 期望运行状态, 是 sync loop 的权威依据。
|
||||
# 手动启动必须置 True, 否则下轮 sync 会把它当"多余实例"停掉。
|
||||
inst.enabled = True
|
||||
inst.running = True
|
||||
db.session.commit()
|
||||
return True
|
||||
|
||||
def stop_instance(self, instance_id):
|
||||
"""停止单个实例。"""
|
||||
"""停止单个实例。
|
||||
|
||||
关键: 必须同时置 enabled=False。enabled 是 sync loop 判断"该实例是否
|
||||
应该在跑"的唯一依据——只停进程不改 enabled, 30 秒内 sync loop 会把它
|
||||
当"死掉的实例"自动重启, 用户的手动停止就失效了。
|
||||
"""
|
||||
from models import Instance
|
||||
from database import db
|
||||
with db.app.app_context():
|
||||
@@ -268,7 +321,8 @@ class InstanceManager:
|
||||
srv = self._instances.pop(inst.name, None)
|
||||
if srv:
|
||||
srv.stop()
|
||||
del self._config_cache[inst.name]
|
||||
self._config_cache.pop(inst.name, None)
|
||||
inst.enabled = False # 手动停止: sync loop 不得自动重启
|
||||
inst.running = False
|
||||
inst.active_connections = 0
|
||||
db.session.commit()
|
||||
|
||||
+21
-3
@@ -373,13 +373,27 @@ class ConnectionHandler:
|
||||
await self._record_stats()
|
||||
|
||||
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:
|
||||
return
|
||||
total_transferred = 0
|
||||
# 读超时短于 config.timeout 时宁可提前踢,不放过慢客户端
|
||||
read_timeout = max(1, int(self.config.timeout))
|
||||
try:
|
||||
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:
|
||||
break
|
||||
|
||||
@@ -455,7 +469,11 @@ class ConnectionHandler:
|
||||
|
||||
# 更新用户流量
|
||||
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:
|
||||
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
|
||||
listen_host = db.Column(db.String(64), default="0.0.0.0")
|
||||
listen_port = db.Column(db.Integer, nullable=False)
|
||||
# 统一超时(秒):覆盖 SOCKS5 握手/认证/请求读取,以及隧道阶段每段 read。
|
||||
# 客户端在任意阶段(含 tunnel)空闲超过此值会被服务端主动断开,避免空连接
|
||||
# 占满 fd 直到 LimitNOFILE 上限。
|
||||
timeout = db.Column(db.Integer, default=30)
|
||||
enabled = db.Column(db.Boolean, default=True)
|
||||
notes = db.Column(db.Text)
|
||||
|
||||
@@ -3,3 +3,4 @@ Flask-SQLAlchemy>=3.1
|
||||
psutil>=5.9
|
||||
bcrypt>=4.0
|
||||
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