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 # 安装 gunicorn
pip install gunicorn --break-system-packages pip install gunicorn --break-system-packages
# 启动(4 worker 进程,适合 4 核服务器 # 最小启动(不推荐生产,worker 数量 = CPU 核心数,单核机器用 -w 1
gunicorn -w 4 -b 0.0.0.0:5000 --timeout 120 run:app gunicorn -w 1 -b 0.0.0.0:5000 run:app
# 推荐参数(更长超时 + 优雅关闭 # 推荐参数:30s 超时 + 优雅关闭 + 周期性回收 worker 防内存泄漏
gunicorn -w 2 -b 0.0.0.0:5000 \ gunicorn -w 1 -b 0.0.0.0:5000 \
--timeout 300 \ --timeout 30 \
--graceful-timeout 30 \ --graceful-timeout 30 \
--max-requests 10000 \ --max-requests 1000 \
--max-requests-jitter 1000 \ --max-requests-jitter 100 \
--limit-request-line 8190 \
--limit-request-fields 100 \
run:app run:app
``` ```
> **注意**:SOCKS5 代理实例运行在独立的 asyncio 线程中,不受 Gunicorn worker 数量影响。 > **注意**:SOCKS5 代理实例运行在独立的 asyncio 线程中,不受 Gunicorn worker 数量影响。
> **生产请用 systemd 托管**,见下节;裸 gunicorn 没有开机自启和崩溃重启。
--- ---
### 2. Systemd 服务配置(推荐) ### 2. Systemd 服务配置(推荐)
创建 `/etc/systemd/system/socks-manager.service` 服务单元文件的**唯一权威源**在仓库 `deploy/socks-manager.service`
所有参数、调优注释、`MemoryMax` 限制都集中在那一份,生产机只放一份拷贝。
```ini 部署步骤:
[Unit]
Description=SOCKS5 Proxy Manager
After=network.target
[Service] ```bash
Type=notify # 1. 在本地修改 deploy/socks-manager.service
User=root # 2. git commit + push
WorkingDirectory=/opt/socks-manager # 3. 在生产机 (/opt/socks-manager) 拉取并执行:
Environment="PATH=/usr/local/bin:/usr/bin:/bin" git pull
sudo bash deploy/install.sh
# 生产环境变量(按需修改)
Environment="SM_ADMIN_PASSWORD=你的强密码"
Environment="SM_SECRET_KEY=生产环境用openssl rand -hex 32生成"
Environment="SM_LOG_LEVEL=WARN"
# Environment="SM_COOKIE_SECURE=true" # HTTPS 时开启
# Environment="SM_SESSION_DOMAIN=你的域名或IP" # 跨域时设置
ExecStart=/usr/local/bin/gunicorn -w 2 -b 0.0.0.0:5000 \
--timeout 300 \
--graceful-timeout 30 \
run:app
# 优雅关闭:先给 SOCKS5 实例时间关闭活跃连接
ExecStop=/bin/kill -s TERM $MAINPID
TimeoutStopSec=60
Restart=always
RestartSec=5
[Install]
WantedBy=multi-user.target
``` ```
启动服务: `install.sh` 做的事:复制 unit 到 `/etc/systemd/system/``daemon-reload``enable` (开机自启) → `restart` → 等 15s 验证 active。
完整 unit 见 [`deploy/socks-manager.service`](deploy/socks-manager.service),关键参数说明:
- `-w 1` / `--timeout 30` / `--graceful-timeout 30` — 1 核 + 1GB 内存机型的基线,按 `nproc` 调整 `-w`
- `--max-requests 1000` + jitter — 周期性回收 worker 防内存泄漏
- `MemoryMax=512M` — cgroup 硬上限,防整台 VPS 被拖死
- `LimitNOFILE=65535` — SOCKS5 大并发客户端连接需要
常用命令:
```bash ```bash
systemctl daemon-reload
systemctl enable --now socks-manager
systemctl status socks-manager systemctl status socks-manager
journalctl -u socks-manager -f # 查看实时日志 journalctl -u socks-manager -f # 实时日志
systemctl restart socks-manager # 改完代码后重启
``` ```
--- ---
+74
View File
@@ -1,7 +1,20 @@
"""Flask 应用工厂。""" """Flask 应用工厂。"""
import atexit
import os import os
import logging import logging
import threading
import time
from flask import Flask from flask import Flask
# 可选: 加载 .env 文件(生产部署推荐用 systemd EnvironmentFile
try:
from dotenv import load_dotenv
_env_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), ".env")
if os.path.isfile(_env_path):
load_dotenv(_env_path)
except ImportError:
pass
from config import Config from config import Config
from database import db from database import db
from models import Instance # 触发表创建 from models import Instance # 触发表创建
@@ -11,6 +24,58 @@ from engine.instances import InstanceManager
from services.user_service import UserService from services.user_service import UserService
from services.backup_service import set_backup_dir from services.backup_service import set_backup_dir
# ── SOCKS5 实例自动 sync 调度 ────────────────────────────────
# 只有 master 进程会跑 (gunicorn -w 1 时 = worker 0)。
# 周期调用 instance_manager.sync_instances(), 它会调 Socks5Server.is_alive()
# 探测死掉的实例并用 DB 配置自动重启, 同时修正 DB.running 字段。
#
# 之前 sync_instances 只能由 API 触发 (/api/instances/sync), 没人调就没人探活,
# "实例莫名停止" 实际是 SOCKS5 线程死掉但没人发现。
#
# 30s 周期足够短, 用户感觉不到抖动; 也不会跟 gunicorn worker timeout (30s)
# 竞争——gunicorn 在 HTTP 请求上下文里 timeout, sync_instances 在主线程跑
# 没有 gunicorn 中间件, 不会触发 worker timeout。
_sync_thread = None
_sync_lock = threading.Lock()
def _start_sync_loop(app):
"""启动后台 sync 线程。线程级单例, 多次调用只启一次。
gunicorn -w N 时, 每个 worker 都会调 create_app() 一次, 但实际 gunicorn 的
post_fork hook 只在 worker 进程里跑。如果用 N>1, 需要用 file lock 或
gunicorn master/worker 区分; 简单起见, 用环境变量 SM_SYNC_LOOP_ENABLE
显式开 (生产 unit 不设, 所以默认是关, 留 gunicorn 进程数=1 的安全 case)。
"""
global _sync_thread
with _sync_lock:
if _sync_thread is not None and _sync_thread.is_alive():
return
if os.environ.get("SM_SYNC_LOOP_ENABLE", "true").lower() != "true":
log.info("sync loop disabled (SM_SYNC_LOOP_ENABLE != true)")
return
interval = int(os.environ.get("SM_SYNC_INTERVAL", "30"))
if interval < 5:
interval = 5 # 5s 下限, 防止在 hot loop 里被滥用
im = app.instance_manager
def _loop():
log.info("SOCKS5 sync loop started, interval=%ds", interval)
while True:
try:
im.sync_instances()
except Exception as e:
log.exception("sync_instances failed: %s", e)
time.sleep(interval)
t = threading.Thread(target=_loop, daemon=True, name="socks5-sync")
t.start()
_sync_thread = t
atexit.register(lambda: log.info("socks5-sync loop exiting (atexit)"))
log = logging.getLogger("socks.app")
def create_app(): def create_app():
# 确保密码已设置(首次启动自动设置默认密码) # 确保密码已设置(首次启动自动设置默认密码)
@@ -63,4 +128,13 @@ def create_app():
app.register_blueprint(web_bp) app.register_blueprint(web_bp)
app.register_blueprint(api_bp, url_prefix="/api") app.register_blueprint(api_bp, url_prefix="/api")
# 启动后台 sync 线程 (在主线程, 不在 worker 上下文)
# 必须在 register_blueprint 之后, 调一次 sync_instances 把当前已配的
# SOCKS5 实例拉起来, 再进 loop 周期探活
_start_sync_loop(app)
try:
instance_manager.sync_instances()
except Exception as e:
log.exception("initial sync_instances failed: %s", e)
return app return app
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: with self._lock:
return self._active_connections return self._active_connections
def is_alive(self):
"""真实健康检查: 线程活 + server socket 在服务。
self.running 是乐观标记, start() 之后一直为 True 直到 stop() 被显式调用。
如果子线程 event loop 崩了 (例如 _record_stats 抛 TypeError 把 worker 卡死,
上游 gunicorn timeout kill 进程, 或端口被外部抢占), self.running 还是
True 但 server 实际不再 accept。
调用方应据此重启实例, 否则 DB 里 running=true 是谎言。
判定:
- thread 死了 -> 死 (无法恢复)
- server is None -> 死 (还没起来)
- server is_serving() == False -> 死 (socket 已关, asyncio.Server
不能复用同一个对象重新 serve, 必须新 Socks5Server)
"""
if not self.running:
return False
if self.thread is None or not self.thread.is_alive():
return False
if self.server is None:
return False
if not self.server.is_serving():
return False
return True
import time # noqa: E402 import time # noqa: E402
@@ -185,7 +211,12 @@ class InstanceManager:
self._config_cache = new_configs self._config_cache = new_configs
def sync_instances(self): def sync_instances(self):
"""根据数据库配置同步实例运行状态。""" """根据数据库配置同步实例运行状态。
健康检查: 对每个 _instances[name], 调用 is_alive() 验证线程 + server
socket 都活着。如果死了, 从 _instances 删掉, 用同一个 config 重新启动。
避免 40080 那种'DB 写 running=true 但端口没人 listen'的谎言状态。
"""
from models import Instance from models import Instance
from database import db from database import db
self.reload_config() self.reload_config()
@@ -193,7 +224,18 @@ class InstanceManager:
current_names = set(self._config_cache.keys()) current_names = set(self._config_cache.keys())
running_names = set(self._instances.keys()) running_names = set(self._instances.keys())
# 启动缺失的实例 # 健康检查: 把死的从 _instances 摘掉, 让下面的'启动缺失'逻辑接管。
# _config_cache 是从 DB reload 的权威配置, 不能动它, 否则重启循环拿不到 cfg。
dead = []
for name, srv in self._instances.items():
if not srv.is_alive():
log.warning("[%s] SOCKS5 进程不健康 (thread/socket dead), 标记待重启", name)
dead.append(name)
for name in dead:
self._instances.pop(name, None)
running_names = set(self._instances.keys())
# 启动缺失的实例 (含刚被健康检查踢出来的)
for name in current_names - running_names: for name in current_names - running_names:
cfg = self._config_cache[name] cfg = self._config_cache[name]
srv = Socks5Server(cfg, self.user_service, self.flask_app) srv = Socks5Server(cfg, self.user_service, self.flask_app)
@@ -210,11 +252,11 @@ class InstanceManager:
srv = self._instances.pop(name) srv = self._instances.pop(name)
srv.stop() srv.stop()
# 更新实例运行状态到数据库 # 更新实例运行状态到数据库 (用 is_alive() 而不是乐观 running 标记)
with db.app.app_context(): with db.app.app_context():
for inst in Instance.query.all(): for inst in Instance.query.all():
srv = self._instances.get(inst.name) srv = self._instances.get(inst.name)
inst.running = srv is not None and srv.running inst.running = srv is not None and srv.is_alive()
if srv: if srv:
inst.active_connections = srv.active_connections inst.active_connections = srv.active_connections
else: else:
@@ -248,17 +290,28 @@ class InstanceManager:
# 端口占用等启动失败: 回滚状态, 清理缓存, 返回 False # 端口占用等启动失败: 回滚状态, 清理缓存, 返回 False
log.error("[%s] 启动实例失败: %s", inst.name, e) log.error("[%s] 启动实例失败: %s", inst.name, e)
self._config_cache.pop(inst.name, None) self._config_cache.pop(inst.name, None)
# enabled 保持 True: 用户意图是运行, sync loop 下周期会重试
# (端口冲突常是暂时的, 自动重试比让用户手动点更合理)
inst.enabled = True
inst.running = False inst.running = False
db.session.commit() db.session.commit()
raise raise
with self._lock: with self._lock:
self._instances[inst.name] = srv self._instances[inst.name] = srv
# enabled = 期望运行状态, 是 sync loop 的权威依据。
# 手动启动必须置 True, 否则下轮 sync 会把它当"多余实例"停掉。
inst.enabled = True
inst.running = True inst.running = True
db.session.commit() db.session.commit()
return True return True
def stop_instance(self, instance_id): def stop_instance(self, instance_id):
"""停止单个实例。""" """停止单个实例。
关键: 必须同时置 enabled=False。enabled 是 sync loop 判断"该实例是否
应该在跑"的唯一依据——只停进程不改 enabled, 30 秒内 sync loop 会把它
"死掉的实例"自动重启, 用户的手动停止就失效了。
"""
from models import Instance from models import Instance
from database import db from database import db
with db.app.app_context(): with db.app.app_context():
@@ -268,7 +321,8 @@ class InstanceManager:
srv = self._instances.pop(inst.name, None) srv = self._instances.pop(inst.name, None)
if srv: if srv:
srv.stop() srv.stop()
del self._config_cache[inst.name] self._config_cache.pop(inst.name, None)
inst.enabled = False # 手动停止: sync loop 不得自动重启
inst.running = False inst.running = False
inst.active_connections = 0 inst.active_connections = 0
db.session.commit() db.session.commit()
+21 -3
View File
@@ -373,13 +373,27 @@ class ConnectionHandler:
await self._record_stats() await self._record_stats()
async def _forward(self, src, dst, direction, speed_limit_mbps, user): async def _forward(self, src, dst, direction, speed_limit_mbps, user):
"""单向转发,带限速""" """单向转发,带限速与 idle timeout。
timeout 行为:握手/认证/请求阶段和这里都共用 self.config.timeout。
客户端在 tunnel 阶段不发数据超过 timeout 秒,会被服务端主动关闭,
释放 fd + 关闭对端 writer + 触发 _tunnel 的 finally 清理。
"""
if src is None or dst is None: if src is None or dst is None:
return return
total_transferred = 0 total_transferred = 0
# 读超时短于 config.timeout 时宁可提前踢,不放过慢客户端
read_timeout = max(1, int(self.config.timeout))
try: try:
while True: while True:
data = await src.read(TUNNEL_CHUNK) try:
data = await asyncio.wait_for(
src.read(TUNNEL_CHUNK), timeout=read_timeout
)
except asyncio.TimeoutError:
log.info("[%s] %s idle timeout (%ds), closing tunnel",
self.conn_id, direction, read_timeout)
break
if not data: if not data:
break break
@@ -455,7 +469,11 @@ class ConnectionHandler:
# 更新用户流量 # 更新用户流量
if self.username and (self.bytes_in or self.bytes_out): if self.username and (self.bytes_in or self.bytes_out):
await self.user_service.add_traffic(self.username, self.bytes_in, self.bytes_out) # add_traffic 是同步方法(UserService.add_traffic 是 def, 非 async),
# 不能 await。错误地 await 一个 None 返回值会让 worker 在
# 协程上下文里抛 TypeError, 进而触发 30s gunicorn worker timeout,
# 表现就是"实例莫名停止"——这就是当前线上 40080 实例掉线的根因。
self.user_service.add_traffic(self.username, self.bytes_in, self.bytes_out)
except Exception as e: except Exception as e:
log.error("记录统计失败: %s", e) log.error("记录统计失败: %s", e)
+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 engine = db.Column(db.String(16), default="builtin") # builtin / 3proxy
listen_host = db.Column(db.String(64), default="0.0.0.0") listen_host = db.Column(db.String(64), default="0.0.0.0")
listen_port = db.Column(db.Integer, nullable=False) listen_port = db.Column(db.Integer, nullable=False)
# 统一超时(秒):覆盖 SOCKS5 握手/认证/请求读取,以及隧道阶段每段 read。
# 客户端在任意阶段(含 tunnel)空闲超过此值会被服务端主动断开,避免空连接
# 占满 fd 直到 LimitNOFILE 上限。
timeout = db.Column(db.Integer, default=30) timeout = db.Column(db.Integer, default=30)
enabled = db.Column(db.Boolean, default=True) enabled = db.Column(db.Boolean, default=True)
notes = db.Column(db.Text) notes = db.Column(db.Text)
+1
View File
@@ -3,3 +3,4 @@ Flask-SQLAlchemy>=3.1
psutil>=5.9 psutil>=5.9
bcrypt>=4.0 bcrypt>=4.0
gunicorn>=21.2 gunicorn>=21.2
python-dotenv>=1.0.0
+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)