7 Commits

Author SHA1 Message Date
Your Name 3232159d29 feat: add one-click deploy script + project optimization
- deploy.sh: one-click deployment (venv + pip deps + .env + gunicorn + systemd + firewall)
- gunicorn_config.py: unified gunicorn config with env var overrides
- .env.example: full env var template (admin, security, port, session, backup, etc.)
- app.py: optional dotenv .env file loading
- requirements.txt: add python-dotenv
2026-09-09 13:41:03 +08:00
cnbugs 47714c9622 fix(socks5): stop_instance 置 enabled=False, 防 sync loop 自动重启用户手动停止的实例
回归 (eba7ae8 引入):
  sync loop 周期调 sync_instances, 它只看 DB 里 enabled=True 的实例。
  但 stop_instance 只停进程 + 设 running=False, 不改 enabled。
  结果: 用户手动停止实例后, 30 秒内 sync loop 会把它当'死掉的实例'
  自动拉回来——手动停止完全失效。

语义修正:
  enabled = 期望运行状态, 是 sync loop 判断'该实例是否应该在跑'的唯一权威。
  - stop_instance: enabled=False (手动停止, sync loop 不得自动重启)
  - start_instance: enabled=True (手动启动, 否则下轮 sync 会把它当多余停掉)
  - start 失败: enabled 保持 True (端口冲突常是暂时的, sync loop 自动重试)

测试:
  + tests/e2e_stop_semantics.py
    stop 后等 2 个 loop 周期 (5s), 断言端口不被自动拉回 + DB.enabled=False;
    start 后断言端口恢复 + DB.enabled=True + DB.running=True
  本地 PASS

全套回归:
  smoke tunnel / e2e lifecycle / e2e health / e2e sync loop / e2e stop semantics
  全部 PASS
2026-08-11 00:07:31 +08:00
cnbugs eba7ae831f fix(socks5): 加后台 sync loop, 周期探活自动救活死 SOCKS5 实例
背景:
  is_alive() 健康探活已就绪, 但 sync_instances 只能由 API 手动触发。
  没人调 API 就没人探活, 死掉的 SOCKS5 实例一直挂着, DB.running 还是
  true, 用户看到的就是'实例莫名停止'。这是最后一块拼图。

变更:
  app.py
    + _start_sync_loop()   后台 daemon 线程, 周期调用 sync_instances
                          (默认 30s, SM_SYNC_INTERVAL 可调, SM_SYNC_LOOP_ENABLE
                          可关)。线程级单例, 多次 create_app 只启一次。
    + create_app 末尾调 _start_sync_loop + 初始 sync_instances, 保证
      gunicorn worker 启动时就把实例拉起来

  engine/instances.py
    M is_alive() docstring 补充判定依据 (thread/server/is_serving)

测试:
  + tests/e2e_sync_loop.py  create_app 启动 loop -> 插 DB 实例 -> loop 拉起
                          -> kill server socket -> loop 自动救活 + 换对象
  本地 3 次连续运行全 PASS, 端口 ~4.5s 被救活

验证:
  - smoke tunnel timeout: PASS
  - e2e lifecycle (userpass + add_traffic): PASS
  - e2e health check (sync_instances 探活): PASS
  - e2e sync loop (后台自动救活): PASS
  - app import: OK
2026-08-10 23:57:56 +08:00
cnbugs 74a68fcbd0 fix(socks5): 加 Socks5Server.is_alive() 健康探活, sync_instances 自动重启死实例
背景:
  Socks5Server.running 是乐观标记, start() 之后到 stop() 之前一直 True。
  但子线程 event loop 崩了/端口被外部抢占/worker 被 gunicorn 杀 等情况,
  self.running 仍为 True, DB 里 inst.running=true 是谎言。
  上次 40080 实例莫名停止就是这种状态。

变更:
  + Socks5Server.is_alive()    thread.is_alive() + server.is_serving()
  M InstanceManager.sync_instances  启动前先探活, 死的从 _instances 摘掉
                                  让'启动缺失'循环用 DB 配置重建它
  M inst.running = srv.is_alive()  不再用乐观的 srv.running

测试:
  + tests/e2e_socks5_lifecycle.py   完整 SOCKS5 userpass 生命周期, 拦住
                                   _record_stats 里 await 同步方法的回归
  + tests/e2e_health_check.py      模拟 Socks5Server 死了, sync_instances
                                   自动重启; 反向验证: 注释掉健康检查就 fail

验证:
  - 正向 (修复在): lifecycle PASS, health PASS, smoke tunnel PASS
  - 反向 (回滚修复): 两个 e2e 都正确 FAIL, 证明测试真的能抓回归
2026-08-10 23:45:29 +08:00
cnbugs 22b9427ca5 fix(socks5): 去掉 _record_stats 里对同步方法 add_traffic 的 await
回归: ee8c0b8 引入的 _forward timeout 修复没改 _record_stats,
但烟雾测试里 FakeUserService.add_traffic 误写成 async def 掩盖了
这个 bug。生产机 15:32:08 命中:
  socks.engine: 记录统计失败: object NoneType can't be used in 'await'
  gunicorn: [CRITICAL] WORKER TIMEOUT (pid:428355)
  gunicorn: Error handling request (no URI read)

add_traffic 是 UserService 里的 def (非 async), 不能 await。错误地
await 一个 None 返回值会让 worker 在协程上下文里抛 TypeError,
被 _record_stats 的 except 吃掉, 但 worker 进入不可服务状态,
触发 30s gunicorn timeout, 表现就是 40080 实例莫名停止。

注意 log_event 是 async, 仍需 await; add_traffic 是 def, 不 await。

修法:
  engine/server.py:472  去掉 await, 同步调用 add_traffic
  tests/smoke_tunnel_timeout.py  FakeUserService.add_traffic 改为
                                def, 模拟真实 UserService 签名,
                                防止烟雾测试再次漏掉此类 bug

线上已临时恢复 40080 (curl POST /api/instances/2/restart),
2026-08-10 23:35:30 +08:00
cnbugs ee8c0b83db fix(socks5): 隧道转发阶段补 idle timeout,防慢客户端耗尽 fd
之前 _forward() 直接 await src.read(TUNNEL_CHUNK),无 timeout。
已通过 SOCKS5 握手的客户端可以一直挂着不发数据,占住 fd 不放,
直到 LimitNOFILE=65535 才被内核拒。

握手/认证/请求阶段已经全部用 asyncio.wait_for + config.timeout,
但同一个 timeout 字段没覆盖 tunnel 阶段,语义不一致。

修复:
  engine/server.py:381-407  _forward 每次 read 包 wait_for, 触发
                          TimeoutError 时 break 走 _tunnel 的 finally
                          清理 (关闭 writer, 取消对端 forward 任务)

验证 (tests/smoke_tunnel_timeout.py):
  - 客户端只握手不进 tunnel, 服务端在 config.timeout 秒后关闭
  - 反向: git stash 掉修复, 客户端永远不被关闭, 服务端报
    'Task was destroyed but it is pending' — 证明修复前后行为差
    异真实存在

models.Instance.timeout 注释: 说明该字段覆盖 tunnel 阶段
2026-08-10 23:28:46 +08:00
cnbugs f948a410b6 ops: 把 systemd unit 纳入版本控制,统一 timeout/限流/资源限制
背景:
  生产机 /etc/systemd/system/socks-manager.service 长期脱管:参数跟
  README 文档不一致(-w 2 vs -w 1), 也没有 -w 4 的孤儿 gunicorn 抢占
  5000 端口、 systemd 启不来服务。已删孤儿, 改用本仓库的 unit 作为
  唯一权威源。

变更:
  + deploy/socks-manager.service  单元文件唯一权威源
  + deploy/install.sh              一键同步: install + daemon-reload +
                                  enable + restart + 验证 active
  M README.md                      部署章节改为引用 deploy/, 删除内联
                                  unit 代码(避免再次漂移)

调优:
  - --timeout 30            网络扫描器半截请求占满 worker, 30s 及时回收
  - --max-requests 1000     周期性回收 worker, 防内存泄漏
  - --limit-request-line    防御大请求行导致 worker 长时间阻塞
  - MemoryMax=512M          cgroup 硬上限, 防单 worker 把整台 1GB VPS
                            拖死再拖到 OOM
  - LimitNOFILE=65535       SOCKS5 大并发客户端连接所需
  - Type=notify + 保留 Restart=always / RestartSec=5
2026-08-10 23:06:28 +08:00
16 changed files with 1369 additions and 52 deletions
+38
View File
@@ -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
+30 -43
View File
@@ -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 # 改完代码后重启
```
---
+74
View File
@@ -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
Executable
+297
View File
@@ -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 "================================================================"
+43
View File
@@ -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
+63
View File
@@ -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
View File
@@ -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
View File
@@ -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)
+44
View File
@@ -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")
+3
View File
@@ -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)
+1
View File
@@ -3,3 +3,4 @@ Flask-SQLAlchemy>=3.1
psutil>=5.9
bcrypt>=4.0
gunicorn>=21.2
python-dotenv>=1.0.0
+130
View File
@@ -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)
+185
View File
@@ -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)
+126
View File
@@ -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)
+116
View File
@@ -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)
+138
View File
@@ -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)