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 筛选,分页查询 | | 详细审计日志 | 按事件/用户/IP 筛选,分页查询 |
| 级联代理 | 连接建立时自动代理至目标(可扩展 upstream) | | 级联代理 | 连接建立时自动代理至目标(可扩展 upstream) |
| 负载均衡 | 多实例配置后按端口轮询 | | 负载均衡 | 多实例配置后按端口轮询 |
| 管理员认证 | 环境变量 `SM_ADMIN_PASSWORD` 保护 | | 管理员认证 | 多管理员账户存数据库(`AdminUser` 表,bcrypt 哈希),支持角色权限 |
| 角色权限 | superadmin / admin / viewer 三级角色,细粒度权限位控制 |
| 管理员管理 | Web/API 增删改、启停、改角色、重置密码(带保护校验) |
| 防爆破 | 自动阻断频繁登录尝试(Fail2Ban 机制) | | 防爆破 | 自动阻断频繁登录尝试(Fail2Ban 机制) |
| 密码安全 | bcrypt 哈希存储 SOCKS5 用户密码(兼容旧明文自动迁移) | | 密码安全 | bcrypt 哈希存储 SOCKS5 用户密码(兼容旧明文自动迁移) |
| 会话管理 | 自定义 Cookie 域名/安全/名称,支持 HTTPS 部署 | | 会话管理 | 自定义 Cookie 域名/安全/名称,支持 HTTPS 部署 |
@@ -66,6 +68,35 @@ python3 run.py
> ⚠️ 首次启动时如果 `SM_ADMIN_PASSWORD` 未设置,会自动写入默认密码 `admin123` 到 `.env` 文件。**生产环境请务必修改此密码。** > ⚠️ 首次启动时如果 `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 # 安装 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 # 改完代码后重启
``` ```
--- ---
+103 -1
View File
@@ -1,7 +1,7 @@
"""REST API v1 — 全部路由。""" """REST API v1 — 全部路由。"""
import os import os
from flask import Blueprint, request, jsonify, current_app 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 database import db
from models import Instance, User, AuditLog from models import Instance, User, AuditLog
from services import user_service, stats_service, backup_service 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 keep = request.json.get("keep", 10) if request.is_json else 10
backup_service.cleanup_old_backups(keep=int(keep)) backup_service.cleanup_old_backups(keep=int(keep))
return jsonify({"ok": True}) 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 应用工厂。""" """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,12 +24,60 @@ 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():
# 确保密码已设置(首次启动自动设置默认密码)
from auth import get_or_set_password
get_or_set_password()
app = Flask(__name__) app = Flask(__name__)
app.config.from_object(Config) app.config.from_object(Config)
app.config["SECRET_KEY"] = os.environ.get( app.config["SECRET_KEY"] = os.environ.get(
@@ -35,9 +96,12 @@ def create_app():
db.init_app(app) db.init_app(app)
app.db = db app.db = db
# 全局服务注册 # 建表 + 种子数据(必须在 app_context 中)
with app.app_context(): with app.app_context():
db.create_all() db.create_all()
# 种子初始管理员
from auth import seed_initial_admin
seed_initial_admin()
# 初始化全局 db.app # 初始化全局 db.app
from database import db as _db from database import db as _db
_db.app = app _db.app = app
@@ -63,4 +127,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
+160 -41
View File
@@ -1,20 +1,29 @@
"""登录认证""" """登录认证 + 权限系统。
v2 版本:后台管理员账户存在数据库(AdminUser 表),支持多账户、角色权限。
首次启动时,如果 AdminUser 表为空,会从 SM_ADMIN_USER / SM_ADMIN_PASSWORD
环境变量(或 .env)创建初始超级管理员。
"""
import os import os
import hashlib
import logging import logging
import secrets import secrets
from datetime import datetime, timezone
from functools import wraps 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") log = logging.getLogger("socks.auth")
bp = Blueprint("auth", __name__) 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") pwd = os.environ.get("SM_ADMIN_PASSWORD")
user = os.environ.get("SM_ADMIN_USER", "admin") 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: try:
for line in open(_ENV): for line in open(_ENV):
line = line.strip() line = line.strip()
@@ -29,37 +38,62 @@ def _get_admin_creds():
return (user, pwd or "") 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(): def get_or_set_password():
"""启动时检查密码,如果为空则设置默认值并打印提示。""" """兼容旧版启动流程(run.py 调用)。
import os 现在做的是:确保有初始管理员存在。
u, p = _get_admin_creds() """
if not p or not p.strip(): admin = seed_initial_admin()
default_pwd = "admin123" if admin:
os.environ["SM_ADMIN_PASSWORD"] = default_pwd return (admin.username, "(见数据库)")
os.environ["SM_ADMIN_USER"] = u # 至少返回一个用户名用于显示
_ENV = os.environ.get("SM_ENV_PATH", os.path.join(os.path.dirname(__file__), ".env")) u, p = _read_env_creds()
try: return (u, p or "")
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)
def get_admin_password_hash(): # ── 当前登录用户 ------------------------------------------------------------
"""获取管理员密码哈希(用于安全比较)。"""
_, p = _get_admin_creds() def current_admin():
return hashlib.sha256(p.encode("utf-8")).hexdigest() """返回当前登录的 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): def login_required(f):
@wraps(f) @wraps(f)
def decorated(*args, **kwargs): 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/"): if request.path.startswith("/api/"):
return {"error": "unauthorized"}, 401 return {"error": "unauthorized"}, 401
return redirect(url_for("auth.login", next=request.url)) return redirect(url_for("auth.login", next=request.url))
@@ -67,32 +101,66 @@ def login_required(f):
return decorated 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"]) @bp.route("/login", methods=["GET", "POST"])
def login(): def login():
if session.get("logged_in"): if session.get("logged_in") and session.get("admin_id"):
return redirect(url_for("web.index")) return redirect(url_for("web.index"))
if request.method == "POST": 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_user = request.form.get("username", "").strip()
form_pwd = request.form.get("password", "") form_pwd = request.form.get("password", "")
# 防爆破检查 # 防爆破
client_ip = request.remote_addr or "unknown" client_ip = request.remote_addr or "unknown"
from services.user_service import UserService _fail2ban_record(client_ip, check_only=True)
us = UserService() if _is_banned(client_ip):
blocked, remaining = us.check_fail2ban(client_ip)
if blocked:
log.warning("防爆破已阻断来自 %s 的登录尝试", client_ip) log.warning("防爆破已阻断来自 %s 的登录尝试", client_ip)
return {"error": "登录过于频繁,请稍后再试"}, 429 return {"error": "登录过于频繁,请稍后再试"}, 429
# 使用常量时间比较防止时序攻击 admin = AdminUser.query.filter_by(username=form_user).first()
if form_user == u and secrets.compare_digest(form_pwd, p): if admin and admin.enabled and admin.check_password(form_pwd):
us.record_auth_attempt(client_ip, success=True) # 登录成功
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.permanent = True
session["logged_in"] = 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")) 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 {"error": "用户名或密码错误"}, 401
return _render_login() return _render_login()
@@ -103,6 +171,57 @@ def logout():
return redirect(url_for("auth.login")) 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(): def _render_login():
return """<!doctype html> return """<!doctype html>
<html lang="zh-CN"><head><meta charset="utf-8"> <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: 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")
+64
View File
@@ -4,6 +4,67 @@ from database import db
import bcrypt 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 实例 ────────────────────────────────────────────────── # ── SOCKS5 实例 ──────────────────────────────────────────────────
class Instance(db.Model): class Instance(db.Model):
__tablename__ = "instances" __tablename__ = "instances"
@@ -13,6 +74,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
+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> <div class="sidebar-section">系统</div>
<a class="nav-link {% if nav.system %}active{% endif %}" href="/system"> <a class="nav-link {% if nav.system %}active{% endif %}" href="/system">
<i class="bi bi-shield-lock"></i> 系统管理</a> <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>
</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 管理面板路由。""" """Web 管理面板路由。"""
from flask import Blueprint, render_template, request, redirect, url_for, flash, current_app 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 database import db
from models import Instance, User from models import Instance, User
from services import stats_service, backup_service from services import stats_service, backup_service
@@ -25,6 +25,8 @@ def inject_nav():
page_name = "活跃连接" page_name = "活跃连接"
elif p.startswith("/logs"): elif p.startswith("/logs"):
page_name = "审计日志" page_name = "审计日志"
elif p.startswith("/system/admins"):
page_name = "管理员账户"
elif p.startswith("/system"): elif p.startswith("/system"):
page_name = "系统管理" page_name = "系统管理"
else: else:
@@ -36,7 +38,8 @@ def inject_nav():
"users": p.startswith("/users"), "users": p.startswith("/users"),
"stats": p.startswith("/stats"), "stats": p.startswith("/stats"),
"logs": p.startswith("/logs"), "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, "page_name": page_name,
"logged_in": flask_session.get("logged_in", False), "logged_in": flask_session.get("logged_in", False),
@@ -307,18 +310,23 @@ def logs():
# ── 系统管理 ─────────────────────────────────────────────────── # ── 系统管理 ───────────────────────────────────────────────────
@bp.route("/system") @bp.route("/system")
@login_required @login_required
@permission_required("system:read")
def system(): def system():
import platform, time import platform, time
from auth import current_admin
backups = backup_service.list_backups() backups = backup_service.list_backups()
admin = current_admin()
return render_template("system.html", backups=backups, return render_template("system.html", backups=backups,
sys_version=platform.python_version(), sys_version=platform.python_version(),
sys_platform=platform.platform(), sys_platform=platform.platform(),
sys_hostname=platform.node(), 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"]) @bp.route("/system/backup", methods=["POST"])
@login_required @login_required
@permission_required("system:write")
def create_backup(): def create_backup():
backup_service.create_backup() backup_service.create_backup()
flash("备份已完成", "success") flash("备份已完成", "success")
@@ -327,6 +335,7 @@ def create_backup():
@bp.route("/system/restore/<int:bid>", methods=["POST"]) @bp.route("/system/restore/<int:bid>", methods=["POST"])
@login_required @login_required
@permission_required("system:write")
def restore_backup(bid): def restore_backup(bid):
result = backup_service.restore_backup(bid) result = backup_service.restore_backup(bid)
if "error" in result: if "error" in result:
@@ -334,3 +343,145 @@ def restore_backup(bid):
else: else:
flash("恢复完成,请刷新页面", "success") flash("恢复完成,请刷新页面", "success")
return redirect(url_for("web.system")) 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"))