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
This commit is contained in:
@@ -1,6 +1,9 @@
|
|||||||
"""Flask 应用工厂。"""
|
"""Flask 应用工厂。"""
|
||||||
|
import atexit
|
||||||
import os
|
import os
|
||||||
import logging
|
import logging
|
||||||
|
import threading
|
||||||
|
import time
|
||||||
from flask import Flask
|
from flask import Flask
|
||||||
from config import Config
|
from config import Config
|
||||||
from database import db
|
from database import db
|
||||||
@@ -11,6 +14,58 @@ from engine.instances import InstanceManager
|
|||||||
from services.user_service import UserService
|
from services.user_service import UserService
|
||||||
from services.backup_service import set_backup_dir
|
from services.backup_service import set_backup_dir
|
||||||
|
|
||||||
|
# ── SOCKS5 实例自动 sync 调度 ────────────────────────────────
|
||||||
|
# 只有 master 进程会跑 (gunicorn -w 1 时 = worker 0)。
|
||||||
|
# 周期调用 instance_manager.sync_instances(), 它会调 Socks5Server.is_alive()
|
||||||
|
# 探测死掉的实例并用 DB 配置自动重启, 同时修正 DB.running 字段。
|
||||||
|
#
|
||||||
|
# 之前 sync_instances 只能由 API 触发 (/api/instances/sync), 没人调就没人探活,
|
||||||
|
# "实例莫名停止" 实际是 SOCKS5 线程死掉但没人发现。
|
||||||
|
#
|
||||||
|
# 30s 周期足够短, 用户感觉不到抖动; 也不会跟 gunicorn worker timeout (30s)
|
||||||
|
# 竞争——gunicorn 在 HTTP 请求上下文里 timeout, sync_instances 在主线程跑
|
||||||
|
# 没有 gunicorn 中间件, 不会触发 worker timeout。
|
||||||
|
_sync_thread = None
|
||||||
|
_sync_lock = threading.Lock()
|
||||||
|
|
||||||
|
|
||||||
|
def _start_sync_loop(app):
|
||||||
|
"""启动后台 sync 线程。线程级单例, 多次调用只启一次。
|
||||||
|
|
||||||
|
gunicorn -w N 时, 每个 worker 都会调 create_app() 一次, 但实际 gunicorn 的
|
||||||
|
post_fork hook 只在 worker 进程里跑。如果用 N>1, 需要用 file lock 或
|
||||||
|
gunicorn master/worker 区分; 简单起见, 用环境变量 SM_SYNC_LOOP_ENABLE
|
||||||
|
显式开 (生产 unit 不设, 所以默认是关, 留 gunicorn 进程数=1 的安全 case)。
|
||||||
|
"""
|
||||||
|
global _sync_thread
|
||||||
|
with _sync_lock:
|
||||||
|
if _sync_thread is not None and _sync_thread.is_alive():
|
||||||
|
return
|
||||||
|
if os.environ.get("SM_SYNC_LOOP_ENABLE", "true").lower() != "true":
|
||||||
|
log.info("sync loop disabled (SM_SYNC_LOOP_ENABLE != true)")
|
||||||
|
return
|
||||||
|
interval = int(os.environ.get("SM_SYNC_INTERVAL", "30"))
|
||||||
|
if interval < 5:
|
||||||
|
interval = 5 # 5s 下限, 防止在 hot loop 里被滥用
|
||||||
|
im = app.instance_manager
|
||||||
|
|
||||||
|
def _loop():
|
||||||
|
log.info("SOCKS5 sync loop started, interval=%ds", interval)
|
||||||
|
while True:
|
||||||
|
try:
|
||||||
|
im.sync_instances()
|
||||||
|
except Exception as e:
|
||||||
|
log.exception("sync_instances failed: %s", e)
|
||||||
|
time.sleep(interval)
|
||||||
|
|
||||||
|
t = threading.Thread(target=_loop, daemon=True, name="socks5-sync")
|
||||||
|
t.start()
|
||||||
|
_sync_thread = t
|
||||||
|
atexit.register(lambda: log.info("socks5-sync loop exiting (atexit)"))
|
||||||
|
|
||||||
|
|
||||||
|
log = logging.getLogger("socks.app")
|
||||||
|
|
||||||
|
|
||||||
def create_app():
|
def create_app():
|
||||||
# 确保密码已设置(首次启动自动设置默认密码)
|
# 确保密码已设置(首次启动自动设置默认密码)
|
||||||
@@ -63,4 +118,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
|
||||||
|
|||||||
+6
-1
@@ -159,6 +159,12 @@ class Socks5Server:
|
|||||||
True 但 server 实际不再 accept。
|
True 但 server 实际不再 accept。
|
||||||
|
|
||||||
调用方应据此重启实例, 否则 DB 里 running=true 是谎言。
|
调用方应据此重启实例, 否则 DB 里 running=true 是谎言。
|
||||||
|
|
||||||
|
判定:
|
||||||
|
- thread 死了 -> 死 (无法恢复)
|
||||||
|
- server is None -> 死 (还没起来)
|
||||||
|
- server is_serving() == False -> 死 (socket 已关, asyncio.Server
|
||||||
|
不能复用同一个对象重新 serve, 必须新 Socks5Server)
|
||||||
"""
|
"""
|
||||||
if not self.running:
|
if not self.running:
|
||||||
return False
|
return False
|
||||||
@@ -166,7 +172,6 @@ class Socks5Server:
|
|||||||
return False
|
return False
|
||||||
if self.server is None:
|
if self.server is None:
|
||||||
return False
|
return False
|
||||||
# server.is_serving() 反映 socket 是否在 listen, 比检查 fd 状态靠谱
|
|
||||||
if not self.server.is_serving():
|
if not self.server.is_serving():
|
||||||
return False
|
return False
|
||||||
return True
|
return True
|
||||||
|
|||||||
@@ -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)
|
||||||
Reference in New Issue
Block a user