8 Commits

Author SHA1 Message Date
Your Name 9890cfc488 feat: add system admin account management with RBAC
- AdminUser model: multi-admin accounts with bcrypt hashing (admin_users table)
- Three roles: superadmin / admin / viewer with granular permission bits
- auth.py: login migrated from env-var single admin to DB-backed accounts;
  seed_initial_admin() auto-creates first superadmin from SM_ADMIN_PASSWORD
- Web UI: /system/admins page (add/edit role/toggle/reset pwd/delete)
  + change-my-password; sidebar entry; permission-guarded routes
- REST API: /api/admins CRUD with protection checks
- Protections: cannot delete/disable self; keep >=1 enabled superadmin;
  disabled accounts fail permission checks immediately
- README: document roles, permission bits, API
2026-09-09 14:33:59 +08:00
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
21 changed files with 2088 additions and 103 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
+62 -44
View File
@@ -39,7 +39,9 @@
| 详细审计日志 | 按事件/用户/IP 筛选,分页查询 |
| 级联代理 | 连接建立时自动代理至目标(可扩展 upstream) |
| 负载均衡 | 多实例配置后按端口轮询 |
| 管理员认证 | 环境变量 `SM_ADMIN_PASSWORD` 保护 |
| 管理员认证 | 多管理员账户存数据库(`AdminUser` 表,bcrypt 哈希),支持角色权限 |
| 角色权限 | superadmin / admin / viewer 三级角色,细粒度权限位控制 |
| 管理员管理 | Web/API 增删改、启停、改角色、重置密码(带保护校验) |
| 防爆破 | 自动阻断频繁登录尝试(Fail2Ban 机制) |
| 密码安全 | bcrypt 哈希存储 SOCKS5 用户密码(兼容旧明文自动迁移) |
| 会话管理 | 自定义 Cookie 域名/安全/名称,支持 HTTPS 部署 |
@@ -66,6 +68,35 @@ python3 run.py
> ⚠️ 首次启动时如果 `SM_ADMIN_PASSWORD` 未设置,会自动写入默认密码 `admin123` 到 `.env` 文件。**生产环境请务必修改此密码。**
## 系统账户管理
Web 后台管理员账户存储在数据库(`admin_users` 表),支持多账户与角色权限,通过 **系统管理 → 管理员账户** 管理(也可用 REST API `/api/admins`)。
### 三级角色
| 角色 | 权限 |
|------|------|
| **superadmin** 超级管理员 | 全部权限,含管理员账户管理(增删改/启停/改角色/重置密码) |
| **admin** 管理员 | 用户/实例/系统管理,可查看管理员账户但不能修改 |
| **viewer** 只读用户 | 仅查看仪表盘/实例/用户/日志,无任何写操作 |
### 角色权限位
每个角色对应一组权限位(`users:read``users:write``instances:read``instances:write``system:read``system:write``admins:read``admins:write``logs:read``stats:read`),可在 `models.py``ROLE_PERMISSIONS` 中扩展。
### 安全保护
- **不能删除/禁用当前登录账户**
- **至少保留一个启用的超级管理员**(删除最后一位 superadmin 会被拒绝)
- 登录密码 bcrypt 哈希存储
- 账户禁用后立即失效(权限校验实时检查)
### REST API
```
GET /api/admins # 列出管理员
POST /api/admins # 创建 {username, password, role}
PUT /api/admins/<id> # 更新 {role?, password?, enabled?}
DELETE /api/admins/<id> # 删除(带保护校验)
```
## 环境变量
| 变量 | 默认值 | 说明 |
@@ -187,66 +218,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 # 改完代码后重启
```
---
+103 -1
View File
@@ -1,7 +1,7 @@
"""REST API v1 — 全部路由。"""
import os
from flask import Blueprint, request, jsonify, current_app
from auth import login_required
from auth import login_required, permission_required
from database import db
from models import Instance, User, AuditLog
from services import user_service, stats_service, backup_service
@@ -279,3 +279,105 @@ def cleanup_backups():
keep = request.json.get("keep", 10) if request.is_json else 10
backup_service.cleanup_old_backups(keep=int(keep))
return jsonify({"ok": True})
# ═══════════════════════════════════════════════════════════════
# 后台管理员账户
# ═══════════════════════════════════════════════════════════════
def _serialize_admin(a):
return {
"id": a.id,
"username": a.username,
"role": a.role,
"role_label": a.role_label,
"enabled": a.enabled,
"last_login_at": a.last_login_at.isoformat() if a.last_login_at else None,
"last_login_ip": a.last_login_ip,
"created_at": a.created_at.isoformat() if a.created_at else None,
}
@bp.route("/admins")
@login_required
@permission_required("admins:read")
def list_admins():
from models import AdminUser
admins = AdminUser.query.order_by(AdminUser.id.asc()).all()
return jsonify([_serialize_admin(a) for a in admins])
@bp.route("/admins", methods=["POST"])
@login_required
@permission_required("admins:write")
def create_admin():
from models import AdminUser
data = request.get_json(silent=True) or {}
username = (data.get("username") or "").strip()
password = data.get("password") or ""
role = data.get("role", "viewer")
if not username or not password:
return jsonify({"error": "用户名和密码不能为空"}), 400
if len(password) < 4:
return jsonify({"error": "密码至少 4 位"}), 400
if role not in AdminUser.ROLE_LABELS:
role = "viewer"
if AdminUser.query.filter_by(username=username).first():
return jsonify({"error": "用户名已存在"}), 409
admin = AdminUser()
admin.username = username
admin.role = role
admin.enabled = True
admin.set_password(password)
db.session.add(admin)
db.session.commit()
return jsonify(_serialize_admin(admin)), 201
@bp.route("/admins/<int:aid>", methods=["DELETE"])
@login_required
@permission_required("admins:write")
def delete_admin_api(aid):
from models import AdminUser
from auth import current_admin
admin = AdminUser.query.get_or_404(aid)
curr = current_admin()
if curr and curr.id == admin.id:
return jsonify({"error": "不能删除当前登录账户"}), 400
if admin.role == "superadmin":
super_count = AdminUser.query.filter_by(role="superadmin", enabled=True).count()
if super_count <= 1:
return jsonify({"error": "至少需要保留一个启用的超级管理员"}), 400
db.session.delete(admin)
db.session.commit()
return jsonify({"ok": True, "deleted": admin.username})
@bp.route("/admins/<int:aid>", methods=["PUT"])
@login_required
@permission_required("admins:write")
def update_admin_api(aid):
from models import AdminUser
admin = AdminUser.query.get_or_404(aid)
data = request.get_json(silent=True) or {}
# 角色更新
role = data.get("role")
if role:
if role not in AdminUser.ROLE_LABELS:
return jsonify({"error": "无效的角色"}), 400
admin.role = role
# 密码更新
password = data.get("password")
if password:
if len(password) < 4:
return jsonify({"error": "密码至少 4 位"}), 400
admin.set_password(password)
# 启停
if "enabled" in data:
enabled = bool(data["enabled"])
from auth import current_admin
curr = current_admin()
if curr and curr.id == admin.id and not enabled:
return jsonify({"error": "不能禁用当前登录账户"}), 400
admin.enabled = enabled
db.session.commit()
return jsonify(_serialize_admin(admin))
+78 -5
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,12 +24,60 @@ 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():
# 确保密码已设置(首次启动自动设置默认密码)
from auth import get_or_set_password
get_or_set_password()
app = Flask(__name__)
app.config.from_object(Config)
app.config["SECRET_KEY"] = os.environ.get(
@@ -35,9 +96,12 @@ def create_app():
db.init_app(app)
app.db = db
# 全局服务注册
# 建表 + 种子数据(必须在 app_context 中)
with app.app_context():
db.create_all()
# 种子初始管理员
from auth import seed_initial_admin
seed_initial_admin()
# 初始化全局 db.app
from database import db as _db
_db.app = app
@@ -63,4 +127,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
+160 -41
View File
@@ -1,20 +1,29 @@
"""登录认证"""
"""登录认证 + 权限系统。
v2 版本:后台管理员账户存在数据库(AdminUser 表),支持多账户、角色权限。
首次启动时,如果 AdminUser 表为空,会从 SM_ADMIN_USER / SM_ADMIN_PASSWORD
环境变量(或 .env)创建初始超级管理员。
"""
import os
import hashlib
import logging
import secrets
from datetime import datetime, timezone
from functools import wraps
from flask import Blueprint, request, session, redirect, url_for, current_app
from flask import Blueprint, request, session, redirect, url_for, current_app, abort
log = logging.getLogger("socks.auth")
bp = Blueprint("auth", __name__)
def _get_admin_creds():
_ENV = os.environ.get("SM_ENV_PATH", os.path.join(os.path.dirname(__file__), ".env"))
# ── 工具函数 ----------------------------------------------------------------
def _read_env_creds():
"""从环境变量或 .env 文件读取初始管理员凭据。"""
_ENV = os.environ.get("SM_ENV_PATH",
os.path.join(os.path.dirname(__file__), ".env"))
pwd = os.environ.get("SM_ADMIN_PASSWORD")
user = os.environ.get("SM_ADMIN_USER", "admin")
if not pwd and os.path.exists(_ENV):
if (not pwd) and os.path.exists(_ENV):
try:
for line in open(_ENV):
line = line.strip()
@@ -29,37 +38,62 @@ def _get_admin_creds():
return (user, pwd or "")
def seed_initial_admin():
"""确保至少存在一个超级管理员。
- 如果 admin_users 表为空:从环境变量/.env 创建初始 superadmin
- 如果表非空:什么都不做
返回创建的 admin 对象或 None。
"""
from models import AdminUser
from database import db
if AdminUser.query.count() > 0:
return None
user, pwd = _read_env_creds()
if not pwd or not pwd.strip():
pwd = "admin123"
log.warning("SM_ADMIN_PASSWORD 未设置,使用默认密码 admin123")
admin = AdminUser()
admin.username = user or "admin"
admin.role = "superadmin"
admin.enabled = True
admin.set_password(pwd)
db.session.add(admin)
db.session.commit()
log.info("已创建初始管理员: %s (role=superadmin)", admin.username)
return admin
def get_or_set_password():
"""启动时检查密码,如果为空则设置默认值并打印提示。"""
import os
u, p = _get_admin_creds()
if not p or not p.strip():
default_pwd = "admin123"
os.environ["SM_ADMIN_PASSWORD"] = default_pwd
os.environ["SM_ADMIN_USER"] = u
_ENV = os.environ.get("SM_ENV_PATH", os.path.join(os.path.dirname(__file__), ".env"))
try:
with open(_ENV, "w") as f:
f.write(f"SM_ADMIN_USER={u}\nSM_ADMIN_PASSWORD={default_pwd}\n")
print(f"⚠️ 密码未设置,已使用默认密码: {default_pwd}")
print(f" 配置文件: {_ENV}")
print(f" 请立即修改密码!")
except Exception as e:
print(f"⚠️ 密码未设置,默认: {default_pwd} (无法写入 .env: {e})")
return (u, default_pwd)
return (u, p)
"""兼容旧版启动流程(run.py 调用)。
现在做的是:确保有初始管理员存在。
"""
admin = seed_initial_admin()
if admin:
return (admin.username, "(见数据库)")
# 至少返回一个用户名用于显示
u, p = _read_env_creds()
return (u, p or "")
def get_admin_password_hash():
"""获取管理员密码哈希(用于安全比较)。"""
_, p = _get_admin_creds()
return hashlib.sha256(p.encode("utf-8")).hexdigest()
# ── 当前登录用户 ------------------------------------------------------------
def current_admin():
"""返回当前登录的 AdminUser 对象;未登录返回 None。"""
from models import AdminUser
uid = session.get("admin_id")
if not uid:
return None
return AdminUser.query.get(uid)
def login_required(f):
@wraps(f)
def decorated(*args, **kwargs):
if not session.get("logged_in"):
if not session.get("logged_in") or not session.get("admin_id"):
if request.path.startswith("/api/"):
return {"error": "unauthorized"}, 401
return redirect(url_for("auth.login", next=request.url))
@@ -67,32 +101,66 @@ def login_required(f):
return decorated
def permission_required(perm):
"""权限装饰器:需要指定权限位才能访问。"""
def decorator(f):
@wraps(f)
def decorated(*args, **kwargs):
admin = current_admin()
if not admin:
if request.path.startswith("/api/"):
return {"error": "unauthorized"}, 401
return redirect(url_for("auth.login", next=request.url))
if not admin.enabled:
session.clear()
return {"error": "账户已被禁用"}, 403
if not admin.has_permission(perm):
if request.path.startswith("/api/"):
return {"error": "forbidden"}, 403
abort(403)
return f(*args, **kwargs)
return decorated
return decorator
# ── 登录 / 登出 -------------------------------------------------------------
@bp.route("/login", methods=["GET", "POST"])
def login():
if session.get("logged_in"):
if session.get("logged_in") and session.get("admin_id"):
return redirect(url_for("web.index"))
if request.method == "POST":
u, p = _get_admin_creds()
from models import AdminUser
from database import db
form_user = request.form.get("username", "").strip()
form_pwd = request.form.get("password", "")
# 防爆破检查
# 防爆破
client_ip = request.remote_addr or "unknown"
from services.user_service import UserService
us = UserService()
blocked, remaining = us.check_fail2ban(client_ip)
if blocked:
_fail2ban_record(client_ip, check_only=True)
if _is_banned(client_ip):
log.warning("防爆破已阻断来自 %s 的登录尝试", client_ip)
return {"error": "登录过于频繁,请稍后再试"}, 429
# 使用常量时间比较防止时序攻击
if form_user == u and secrets.compare_digest(form_pwd, p):
us.record_auth_attempt(client_ip, success=True)
admin = AdminUser.query.filter_by(username=form_user).first()
if admin and admin.enabled and admin.check_password(form_pwd):
# 登录成功
admin.last_login_at = datetime.now(timezone.utc)
admin.last_login_ip = client_ip
db.session.commit()
_reset_fail2ban(client_ip)
session.permanent = True
session["logged_in"] = True
session["user"] = u
session["admin_id"] = admin.id
session["user"] = admin.username
session["role"] = admin.role
return redirect(request.args.get("next") or url_for("web.index"))
us.record_auth_attempt(client_ip, success=False)
# 登录失败
_record_fail2ban(client_ip)
return {"error": "用户名或密码错误"}, 401
return _render_login()
@@ -103,6 +171,57 @@ def logout():
return redirect(url_for("auth.login"))
# ── 防爆破(进程内简易实现) ------------------------------------------------
# 注意:多 worker 下每个进程独立,生产环境建议用 Redis 替代。
_fail_attempts = {} # {ip: [timestamps...]}
_ban_until = {} # {ip: timestamp}
_BAN_WINDOW = 600 # 10 分钟窗口
_BAN_MAX = 10 # 最多 10 次失败
_BAN_DURATION = 900 # 封禁 15 分钟
def _is_banned(ip):
now = _now()
if ip in _ban_until and _ban_until[ip] > now:
return True
if ip in _ban_until:
del _ban_until[ip]
return False
def _fail2ban_record(ip, check_only=False):
"""检查 + 清理过期记录(兼容旧版 UserService.check_fail2ban 调用方式)。"""
now = _now()
if ip in _fail_attempts:
_fail_attempts[ip] = [t for t in _fail_attempts[ip] if now - t < _BAN_WINDOW]
return _is_banned(ip), _BAN_MAX - len(_fail_attempts.get(ip, []))
def _record_fail2ban(ip):
now = _now()
if ip not in _fail_attempts:
_fail_attempts[ip] = []
_fail_attempts[ip].append(now)
_fail_attempts[ip] = [t for t in _fail_attempts[ip] if now - t < _BAN_WINDOW]
if len(_fail_attempts[ip]) >= _BAN_MAX:
_ban_until[ip] = now + _BAN_DURATION
log.warning("IP %s 登录失败 %d 次,已封禁 %d",
ip, _BAN_MAX, _BAN_DURATION)
def _reset_fail2ban(ip):
_fail_attempts.pop(ip, None)
_ban_until.pop(ip, None)
def _now():
import time
return time.time()
# ── 登录页 -----------------------------------------------------------------
def _render_login():
return """<!doctype html>
<html lang="zh-CN"><head><meta charset="utf-8">
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")
+64
View File
@@ -4,6 +4,67 @@ from database import db
import bcrypt
# ── 后台管理员 ──────────────────────────────────────────────────
class AdminUser(db.Model):
"""Web 后台管理员账户(存数据库,支持多账户、角色权限)。"""
__tablename__ = "admin_users"
id = db.Column(db.Integer, primary_key=True)
username = db.Column(db.String(64), nullable=False, unique=True, index=True)
password = db.Column(db.String(128), nullable=False) # bcrypt 哈希
role = db.Column(db.String(16), default="viewer") # superadmin | admin | viewer
enabled = db.Column(db.Boolean, default=True)
last_login_at = db.Column(db.DateTime, nullable=True)
last_login_ip = db.Column(db.String(64), nullable=True)
created_at = db.Column(db.DateTime, default=lambda: datetime.now(timezone.utc))
updated_at = db.Column(db.DateTime, default=lambda: datetime.now(timezone.utc),
onupdate=lambda: datetime.now(timezone.utc))
ROLE_LABELS = {
"superadmin": "超级管理员",
"admin": "管理员",
"viewer": "只读用户",
}
# 角色权限位(可扩展)
ROLE_PERMISSIONS = {
"superadmin": {"users:read", "users:write", "instances:read", "instances:write",
"system:read", "system:write", "admins:read", "admins:write",
"logs:read", "stats:read"},
"admin": {"users:read", "users:write", "instances:read", "instances:write",
"system:read", "system:write", "admins:read",
"logs:read", "stats:read"},
"viewer": {"users:read", "instances:read", "system:read",
"logs:read", "stats:read"},
}
def set_password(self, password):
self.password = bcrypt.hashpw(
password.encode("utf-8"), bcrypt.gensalt()
).decode("utf-8")
def check_password(self, password):
if not self.password:
return False
try:
return bcrypt.checkpw(
password.encode("utf-8"), self.password.encode("utf-8")
)
except Exception:
return False
@property
def role_label(self):
return self.ROLE_LABELS.get(self.role, self.role)
def has_permission(self, perm):
perms = self.ROLE_PERMISSIONS.get(self.role, set())
return perm in perms
def __repr__(self):
return f"<AdminUser {self.username} ({self.role})>"
# ── SOCKS5 实例 ──────────────────────────────────────────────────
class Instance(db.Model):
__tablename__ = "instances"
@@ -13,6 +74,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
+203
View File
@@ -0,0 +1,203 @@
{% extends "base.html" %}
{% block title %}管理员账户 · SOCKS Manager{% endblock %}
{% block content %}
<div class="page-header">
<h5><i class="bi bi-person-gear"></i> 管理员账户</h5>
</div>
{% set can_write = current_admin and current_admin.has_permission('admins:write') %}
<!-- 新增管理员 -->
<div class="card mb-3">
<div class="card-header"><span><i class="bi bi-person-plus"></i> 新增管理员</span></div>
<div class="card-body">
<form method="post" action="/system/admins/add" class="row g-2 align-items-end">
<div class="col-md-3">
<label class="form-label small">用户名</label>
<input name="username" class="form-control" placeholder="如: zhangsan" required>
</div>
<div class="col-md-3">
<label class="form-label small">密码</label>
<input name="password" type="password" class="form-control" placeholder="至少 4 位" required>
</div>
<div class="col-md-3">
<label class="form-label small">角色</label>
<select name="role" class="form-select">
{% for r, label in [('superadmin','超级管理员'), ('admin','管理员'), ('viewer','只读用户')] %}
<option value="{{ r }}">{{ label }}</option>
{% endfor %}
</select>
</div>
<div class="col-md-3">
<button type="submit" class="btn btn-primary w-100">
<i class="bi bi-plus-lg"></i> 创建</button>
</div>
</form>
</div>
</div>
<!-- 角色说明 -->
<div class="card mb-3">
<div class="card-header"><span><i class="bi bi-info-circle"></i> 角色权限说明</span></div>
<div class="card-body">
<div class="row g-3">
<div class="col-md-4">
<div class="card" style="background:var(--bg-tertiary)">
<div class="card-body py-2">
<span class="badge badge-danger">超级管理员</span>
<span class="text-muted small">全部权限,含管理员管理</span>
</div>
</div>
</div>
<div class="col-md-4">
<div class="card" style="background:var(--bg-tertiary)">
<div class="card-body py-2">
<span class="badge badge-info">管理员</span>
<span class="text-muted small">用户/实例/系统管理,可查看管理员</span>
</div>
</div>
</div>
<div class="col-md-4">
<div class="card" style="background:var(--bg-tertiary)">
<div class="card-body py-2">
<span class="badge badge-secondary">只读用户</span>
<span class="text-muted small">仅查看,不可修改</span>
</div>
</div>
</div>
</div>
</div>
</div>
<!-- 管理员列表 -->
<div class="card">
<div class="card-header"><span><i class="bi bi-people"></i> 账户列表</span></div>
<div class="card-body p-0">
<table class="table table-dark table-sm">
<thead><tr>
<th>用户名</th><th>角色</th><th>状态</th><th>最近登录</th><th>创建时间</th>
{% if can_write %}<th>操作</th>{% endif %}
</tr></thead>
<tbody>
{% for a in admins %}
<tr>
<td>
{{ a.username }}
{% if current_admin and current_admin.id == a.id %}
<span class="badge badge-info">当前</span>
{% endif %}
</td>
<td>
{% if can_write %}
<form method="post" action="/system/admins/{{ a.id }}/role" class="d-inline">
<select name="role" class="form-select d-inline-block" style="width:auto;padding:.15rem .4rem;font-size:12px"
onchange="this.form.submit()">
{% for r, label in [('superadmin','超级管理员'), ('admin','管理员'), ('viewer','只读用户')] %}
<option value="{{ r }}" {% if a.role == r %}selected{% endif %}>{{ label }}</option>
{% endfor %}
</select>
</form>
{% else %}
<span class="badge {{ 'badge-danger' if a.role=='superadmin' else 'badge-info' if a.role=='admin' else 'badge-secondary' }}">{{ a.role_label }}</span>
{% endif %}
</td>
<td>
{% if a.enabled %}
<span class="badge badge-success">启用</span>
{% else %}
<span class="badge badge-danger">禁用</span>
{% endif %}
</td>
<td class="text-muted">
{% if a.last_login_at %}
{{ a.last_login_at.strftime('%Y-%m-%d %H:%M') }}<br>
<small>{{ a.last_login_ip or '-' }}</small>
{% else %}-{% endif %}
</td>
<td class="text-muted">{{ a.created_at.strftime('%Y-%m-%d') if a.created_at else '-' }}</td>
{% if can_write %}
<td>
{% if not (current_admin and current_admin.id == a.id) %}
<form method="post" action="/system/admins/{{ a.id }}/toggle" class="d-inline">
<button class="btn btn-sm {{ 'btn-outline-warning' if a.enabled else 'btn-outline-success' }}"
title="{{ '禁用' if a.enabled else '启用' }}">
<i class="bi {{ 'bi-pause-circle' if a.enabled else 'bi-play-circle' }}"></i></button>
</form>
<!-- 重置密码(模态框) -->
<button class="btn btn-sm btn-outline-primary" onclick="openPwdModal({{ a.id }}, '{{ a.username }}')">
<i class="bi bi-key"></i> 重置密码</button>
<form method="post" action="/system/admins/{{ a.id }}/delete" class="d-inline"
onsubmit="return confirm('确认删除管理员 {{ a.username }}')">
<button class="btn btn-sm btn-outline-danger"><i class="bi bi-trash"></i></button>
</form>
{% else %}
<span class="text-muted small"></span>
{% endif %}
</td>
{% endif %}
</tr>
{% endfor %}
</tbody>
</table>
</div>
</div>
<!-- 修改当前密码 -->
<div class="card mt-3">
<div class="card-header"><span><i class="bi bi-shield-check"></i> 修改我的密码</span></div>
<div class="card-body">
<form method="post" action="/system/change-password" class="row g-2 align-items-end">
<div class="col-md-3">
<label class="form-label small">原密码</label>
<input name="old_password" type="password" class="form-control" required>
</div>
<div class="col-md-3">
<label class="form-label small">新密码</label>
<input name="new_password" type="password" class="form-control" required>
</div>
<div class="col-md-3">
<label class="form-label small">确认新密码</label>
<input name="confirm_password" type="password" class="form-control" required>
</div>
<div class="col-md-3">
<button type="submit" class="btn btn-outline-primary w-100">
<i class="bi bi-check-lg"></i> 修改密码</button>
</div>
</form>
</div>
</div>
<!-- 重置密码模态框 -->
<div class="modal fade" id="pwdModal" tabindex="-1">
<div class="modal-dialog">
<div class="modal-content">
<form method="post" action="" id="pwdForm">
<div class="modal-header">
<h5 class="modal-title">重置密码</h5>
<button type="button" class="btn-close" data-bs-dismiss="modal"></button>
</div>
<div class="modal-body">
<p class="text-muted small" id="pwdTarget"></p>
<label class="form-label small">新密码(至少 4 位)</label>
<input name="password" type="password" class="form-control" required>
</div>
<div class="modal-footer">
<button type="button" class="btn btn-outline-secondary" data-bs-dismiss="modal">取消</button>
<button type="submit" class="btn btn-primary">确认重置</button>
</div>
</form>
</div>
</div>
</div>
{% endblock %}
{% block extra_script %}
<script>
function openPwdModal(id, username) {
document.getElementById('pwdForm').action = '/system/admins/' + id + '/password';
document.getElementById('pwdTarget').textContent = '账户: ' + username;
new bootstrap.Modal(document.getElementById('pwdModal')).show();
}
</script>
{% endblock %}
+2
View File
@@ -217,6 +217,8 @@ canvas { max-width:100%; }
<div class="sidebar-section">系统</div>
<a class="nav-link {% if nav.system %}active{% endif %}" href="/system">
<i class="bi bi-shield-lock"></i> 系统管理</a>
<a class="nav-link {% if nav.admins %}active{% endif %}" href="/system/admins">
<i class="bi bi-person-gear"></i> 管理员账户</a>
</div>
</div>
+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)
+154 -3
View File
@@ -1,6 +1,6 @@
"""Web 管理面板路由。"""
from flask import Blueprint, render_template, request, redirect, url_for, flash, current_app
from auth import login_required
from auth import login_required, permission_required
from database import db
from models import Instance, User
from services import stats_service, backup_service
@@ -25,6 +25,8 @@ def inject_nav():
page_name = "活跃连接"
elif p.startswith("/logs"):
page_name = "审计日志"
elif p.startswith("/system/admins"):
page_name = "管理员账户"
elif p.startswith("/system"):
page_name = "系统管理"
else:
@@ -36,7 +38,8 @@ def inject_nav():
"users": p.startswith("/users"),
"stats": p.startswith("/stats"),
"logs": p.startswith("/logs"),
"system": p.startswith("/system"),
"system": p.startswith("/system") and not p.startswith("/system/admins"),
"admins": p.startswith("/system/admins"),
},
"page_name": page_name,
"logged_in": flask_session.get("logged_in", False),
@@ -307,18 +310,23 @@ def logs():
# ── 系统管理 ───────────────────────────────────────────────────
@bp.route("/system")
@login_required
@permission_required("system:read")
def system():
import platform, time
from auth import current_admin
backups = backup_service.list_backups()
admin = current_admin()
return render_template("system.html", backups=backups,
sys_version=platform.python_version(),
sys_platform=platform.platform(),
sys_hostname=platform.node(),
uptime=time.strftime("%Y-%m-%d %H:%M:%S", time.localtime()))
uptime=time.strftime("%Y-%m-%d %H:%M:%S", time.localtime()),
current_admin=admin)
@bp.route("/system/backup", methods=["POST"])
@login_required
@permission_required("system:write")
def create_backup():
backup_service.create_backup()
flash("备份已完成", "success")
@@ -327,6 +335,7 @@ def create_backup():
@bp.route("/system/restore/<int:bid>", methods=["POST"])
@login_required
@permission_required("system:write")
def restore_backup(bid):
result = backup_service.restore_backup(bid)
if "error" in result:
@@ -334,3 +343,145 @@ def restore_backup(bid):
else:
flash("恢复完成,请刷新页面", "success")
return redirect(url_for("web.system"))
# ── 管理员账户管理 ────────────────────────────────────────────────
@bp.route("/system/admins")
@login_required
@permission_required("admins:read")
def admins():
from models import AdminUser
from auth import current_admin
admin_list = AdminUser.query.order_by(AdminUser.id.asc()).all()
return render_template("admins.html", admins=admin_list,
current_admin=current_admin())
@bp.route("/system/admins/add", methods=["POST"])
@login_required
@permission_required("admins:write")
def add_admin():
from models import AdminUser
username = request.form.get("username", "").strip()
password = request.form.get("password", "")
role = request.form.get("role", "viewer")
if not username or not password:
flash("用户名和密码不能为空", "danger")
return redirect(url_for("web.admins"))
if role not in AdminUser.ROLE_LABELS:
role = "viewer"
if AdminUser.query.filter_by(username=username).first():
flash("用户名已存在", "danger")
return redirect(url_for("web.admins"))
if len(password) < 4:
flash("密码至少 4 位", "danger")
return redirect(url_for("web.admins"))
admin = AdminUser()
admin.username = username
admin.role = role
admin.enabled = True
admin.set_password(password)
db.session.add(admin)
db.session.commit()
flash(f"已创建管理员 {username}", "success")
return redirect(url_for("web.admins"))
@bp.route("/system/admins/<int:aid>/toggle", methods=["POST"])
@login_required
@permission_required("admins:write")
def toggle_admin(aid):
from models import AdminUser
from auth import current_admin
admin = AdminUser.query.get_or_404(aid)
# 不能禁用自己
curr = current_admin()
if curr and curr.id == admin.id:
flash("不能禁用当前登录账户", "danger")
return redirect(url_for("web.admins"))
admin.enabled = not admin.enabled
db.session.commit()
flash(f"{'已启用' if admin.enabled else '已禁用'} {admin.username}", "info")
return redirect(url_for("web.admins"))
@bp.route("/system/admins/<int:aid>/role", methods=["POST"])
@login_required
@permission_required("admins:write")
def change_admin_role(aid):
from models import AdminUser
admin = AdminUser.query.get_or_404(aid)
role = request.form.get("role", "viewer")
if role not in AdminUser.ROLE_LABELS:
flash("无效的角色", "danger")
return redirect(url_for("web.admins"))
admin.role = role
db.session.commit()
flash(f"{admin.username} 角色已更新为 {admin.role_label}", "success")
return redirect(url_for("web.admins"))
@bp.route("/system/admins/<int:aid>/password", methods=["POST"])
@login_required
@permission_required("admins:write")
def reset_admin_password(aid):
from models import AdminUser
password = request.form.get("password", "")
if not password or len(password) < 4:
flash("密码至少 4 位", "danger")
return redirect(url_for("web.admins"))
admin = AdminUser.query.get_or_404(aid)
admin.set_password(password)
db.session.commit()
flash(f"{admin.username} 密码已重置", "success")
return redirect(url_for("web.admins"))
@bp.route("/system/admins/<int:aid>/delete", methods=["POST"])
@login_required
@permission_required("admins:write")
def delete_admin(aid):
from models import AdminUser
from auth import current_admin
admin = AdminUser.query.get_or_404(aid)
curr = current_admin()
if curr and curr.id == admin.id:
flash("不能删除当前登录账户", "danger")
return redirect(url_for("web.admins"))
# 至少保留一个 superadmin
if admin.role == "superadmin":
super_count = AdminUser.query.filter_by(role="superadmin", enabled=True).count()
if super_count <= 1:
flash("至少需要保留一个启用的超级管理员", "danger")
return redirect(url_for("web.admins"))
db.session.delete(admin)
db.session.commit()
flash(f"已删除管理员 {admin.username}", "warning")
return redirect(url_for("web.admins"))
@bp.route("/system/change-password", methods=["POST"])
@login_required
def change_my_password():
"""当前登录用户修改自己的密码。"""
from models import AdminUser
from auth import current_admin
old_pwd = request.form.get("old_password", "")
new_pwd = request.form.get("new_password", "")
confirm_pwd = request.form.get("confirm_password", "")
admin = current_admin()
if not admin:
return redirect(url_for("auth.logout"))
if not admin.check_password(old_pwd):
flash("原密码错误", "danger")
return redirect(url_for("web.system"))
if not new_pwd or len(new_pwd) < 4:
flash("新密码至少 4 位", "danger")
return redirect(url_for("web.system"))
if new_pwd != confirm_pwd:
flash("两次输入的新密码不一致", "danger")
return redirect(url_for("web.system"))
admin.set_password(new_pwd)
db.session.commit()
flash("密码修改成功", "success")
return redirect(url_for("web.system"))