feat: 插件管理 — 导入安装 + 统计面板
新增 API: - GET /api/plugins/stats — 统计 (总计/已启用/已停用/安装失败) - POST /api/plugins/import — ZIP上传 或 Git clone 安装 - 自动 pip install -r requirements.txt (子进程隔离) - 自动加载插件 (isolated 模式, 崩溃不影响框架) - 安装失败记录到 _IMPORT_FAILS 前端: - 横版统计卡片 (总数/启用/停用/失败) - ZIP 上传按钮 + Git URL 输入 - 安装状态实时反馈 安全: - 所有端点 panel_auth 保护 - 子进程安装 (pip/git), 超时120s - 仅接受 .zip 文件 - 目录冲突检测 Co-Authored-By: Claude <noreply@anthropic.com>
This commit is contained in:
@@ -1,53 +1,244 @@
|
||||
#!/usr/bin/env python3
|
||||
# -*- coding: utf-8 -*-
|
||||
import logging
|
||||
import logging, os, shutil, subprocess, tempfile, zipfile, asyncio, json
|
||||
from pathlib import Path
|
||||
from aiohttp import web
|
||||
from ..utils.auth import panel_auth
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
def setup_routes(app, prefix=''):
|
||||
app.router.add_get(f'{prefix}/api/plugins', panel_auth(list_plugins))
|
||||
app.router.add_get(f'{prefix}/api/plugin-pages', panel_auth(list_plugin_web_pages))
|
||||
app.router.add_post(f'{prefix}/api/plugins/{{name}}/{{action}}', panel_auth(manage_plugin))
|
||||
app.router.add_get(f'{prefix}/api/plugins/{{name}}/perms', panel_auth(get_perms))
|
||||
app.router.add_post(f'{prefix}/api/plugins/{{name}}/perms', panel_auth(set_perms))
|
||||
PLUGINS_DIR = Path("plugins")
|
||||
_IMPORT_FAILS: dict[str, str] = {} # plugin_name → error message
|
||||
|
||||
|
||||
def _run_in_subprocess(cmd: list, cwd: str, timeout: int = 120) -> tuple[bool, str]:
|
||||
"""在子进程中执行命令,不阻塞主事件循环"""
|
||||
try:
|
||||
r = subprocess.run(cmd, cwd=cwd, capture_output=True, text=True, timeout=timeout)
|
||||
return r.returncode == 0, r.stdout[-500:] + "\n" + r.stderr[-500:]
|
||||
except subprocess.TimeoutExpired:
|
||||
return False, "安装超时 (>120s)"
|
||||
except Exception as e:
|
||||
return False, str(e)
|
||||
|
||||
|
||||
def _install_plugin_deps(plugin_dir: str) -> tuple[bool, str]:
|
||||
"""安装插件依赖 (requirements.txt)"""
|
||||
req_file = os.path.join(plugin_dir, "requirements.txt")
|
||||
if not os.path.exists(req_file):
|
||||
return True, ""
|
||||
return _run_in_subprocess(
|
||||
["pip", "install", "-r", "requirements.txt", "--break-system-packages"],
|
||||
plugin_dir,
|
||||
)
|
||||
|
||||
|
||||
async def _async_install_deps(plugin_dir: str) -> tuple[bool, str]:
|
||||
"""异步包装器 — 在线程池中运行 pip install"""
|
||||
loop = asyncio.get_event_loop()
|
||||
return await loop.run_in_executor(None, _install_plugin_deps, plugin_dir)
|
||||
|
||||
|
||||
def setup_routes(app, prefix=""):
|
||||
app.router.add_get(f"{prefix}/api/plugins", panel_auth(list_plugins))
|
||||
app.router.add_get(f"{prefix}/api/plugin-pages", panel_auth(list_plugin_web_pages))
|
||||
app.router.add_get(f"{prefix}/api/plugins/stats", panel_auth(plugin_stats))
|
||||
app.router.add_post(f"{prefix}/api/plugins/{{name}}/{{action}}", panel_auth(manage_plugin))
|
||||
app.router.add_post(f"{prefix}/api/plugins/import", panel_auth(import_plugin))
|
||||
app.router.add_get(f"{prefix}/api/plugins/{{name}}/perms", panel_auth(get_perms))
|
||||
app.router.add_post(f"{prefix}/api/plugins/{{name}}/perms", panel_auth(set_perms))
|
||||
|
||||
|
||||
# ── 统计 ──
|
||||
|
||||
async def plugin_stats(req):
|
||||
sm = req.app.get("service_manager")
|
||||
ps = sm.get_service("plugin") if sm else None
|
||||
total = len(ps.plugin_info) if ps else 0
|
||||
loaded = len(ps.plugins) if ps else 0
|
||||
disabled = total - loaded
|
||||
failed = len(_IMPORT_FAILS)
|
||||
return web.json_response({
|
||||
"total": total,
|
||||
"loaded": loaded,
|
||||
"disabled": max(0, disabled),
|
||||
"failed": failed,
|
||||
"fail_details": _IMPORT_FAILS,
|
||||
})
|
||||
|
||||
|
||||
# ── 导入 ──
|
||||
|
||||
async def import_plugin(req):
|
||||
"""导入插件: multipart zip 或 JSON {url: git地址}"""
|
||||
sm = req.app.get("service_manager")
|
||||
ps = sm.get_service("plugin") if sm else None
|
||||
if not ps:
|
||||
return web.json_response({"ok": False, "error": "Plugin Service 未就绪"}, status=503)
|
||||
|
||||
content_type = req.content_type or ""
|
||||
|
||||
# ── Git clone 方式 ──
|
||||
if "application/json" in content_type:
|
||||
try:
|
||||
data = await req.json()
|
||||
git_url = data.get("url", "").strip()
|
||||
if not git_url:
|
||||
return web.json_response({"ok": False, "error": "Git URL 不能为空"}, status=400)
|
||||
|
||||
# 提取 repo 名作为插件名
|
||||
name = git_url.rstrip("/").split("/")[-1]
|
||||
if name.endswith(".git"):
|
||||
name = name[:-4]
|
||||
target = PLUGINS_DIR / name
|
||||
if target.exists():
|
||||
return web.json_response({"ok": False, "error": f"插件目录已存在: {name}"}, status=409)
|
||||
|
||||
ok, out = _run_in_subprocess(
|
||||
["git", "clone", git_url, str(target)], str(PLUGINS_DIR)
|
||||
)
|
||||
if not ok:
|
||||
_IMPORT_FAILS[name] = f"git clone 失败: {out[-200:]}"
|
||||
return web.json_response({"ok": False, "error": "Git clone 失败", "detail": out[-300:]}, status=500)
|
||||
|
||||
# 安装依赖
|
||||
dep_ok, dep_out = await _async_install_deps(str(target))
|
||||
if not dep_ok:
|
||||
logger.warning(f"依赖安装部分失败 {name}: {dep_out[-200:]}")
|
||||
|
||||
# 加载插件
|
||||
loaded = await _safe_load_plugin(ps, name)
|
||||
_IMPORT_FAILS.pop(name, None)
|
||||
return web.json_response({
|
||||
"ok": True, "name": name, "loaded": loaded,
|
||||
"dep_ok": dep_ok, "detail": dep_out[-200:] if not dep_ok else "",
|
||||
})
|
||||
except Exception as e:
|
||||
return web.json_response({"ok": False, "error": str(e)}, status=500)
|
||||
|
||||
# ── ZIP 上传方式 ──
|
||||
try:
|
||||
reader = await req.multipart()
|
||||
part = await reader.next()
|
||||
while part:
|
||||
if part.name == "file" and part.filename:
|
||||
fname = part.filename
|
||||
if not fname.endswith(".zip"):
|
||||
return web.json_response({"ok": False, "error": "仅支持 .zip 文件"}, status=400)
|
||||
|
||||
# 保存到临时文件
|
||||
tmp = tempfile.NamedTemporaryFile(suffix=".zip", delete=False)
|
||||
while True:
|
||||
chunk = await part.read_chunk(65536)
|
||||
if not chunk:
|
||||
break
|
||||
tmp.write(chunk)
|
||||
tmp.close()
|
||||
|
||||
# 提取 → plugins/
|
||||
plugin_name = fname[:-4] # 去掉 .zip
|
||||
target = PLUGINS_DIR / plugin_name
|
||||
if target.exists():
|
||||
os.unlink(tmp.name)
|
||||
return web.json_response({"ok": False, "error": f"插件目录已存在: {plugin_name}"}, status=409)
|
||||
|
||||
try:
|
||||
with zipfile.ZipFile(tmp.name, "r") as zf:
|
||||
# 检查 zip 结构 — 如果顶层只有一个目录,用它
|
||||
members = zf.namelist()
|
||||
top_dirs = set()
|
||||
for m in members:
|
||||
if "/" in m:
|
||||
top_dirs.add(m.split("/")[0])
|
||||
if len(top_dirs) == 1 and all(m.startswith(list(top_dirs)[0]) for m in members if m):
|
||||
# 有统一顶层目录
|
||||
zf.extractall(PLUGINS_DIR)
|
||||
if list(top_dirs)[0] != plugin_name:
|
||||
src = PLUGINS_DIR / list(top_dirs)[0]
|
||||
src.rename(target)
|
||||
else:
|
||||
target.mkdir(parents=True)
|
||||
zf.extractall(target)
|
||||
except zipfile.BadZipFile:
|
||||
os.unlink(tmp.name)
|
||||
return web.json_response({"ok": False, "error": "无效的 ZIP 文件"}, status=400)
|
||||
finally:
|
||||
os.unlink(tmp.name)
|
||||
|
||||
# 安装依赖
|
||||
dep_ok, dep_out = await _async_install_deps(str(target))
|
||||
|
||||
# 加载插件
|
||||
loaded = await _safe_load_plugin(ps, plugin_name)
|
||||
if loaded:
|
||||
_IMPORT_FAILS.pop(plugin_name, None)
|
||||
return web.json_response({
|
||||
"ok": True, "name": plugin_name, "loaded": loaded,
|
||||
"dep_ok": dep_ok, "detail": dep_out[-200:] if not dep_ok else "",
|
||||
})
|
||||
part = await reader.next()
|
||||
return web.json_response({"ok": False, "error": "未找到上传文件 (字段名: file)"}, status=400)
|
||||
except Exception as e:
|
||||
logger.error(f"导入插件失败: {e}")
|
||||
return web.json_response({"ok": False, "error": str(e)}, status=500)
|
||||
|
||||
|
||||
async def _safe_load_plugin(ps, name: str) -> bool:
|
||||
"""安全加载插件 — 先隔离模式尝试,成功后再切换到正常模式"""
|
||||
try:
|
||||
# 先尝试隔离加载 (子进程,崩溃不影响框架)
|
||||
from services.process_isolated import IsolatedPlugin
|
||||
plugin_path = PLUGINS_DIR / name / "__init__.py"
|
||||
if not plugin_path.exists():
|
||||
_IMPORT_FAILS[name] = "缺少 __init__.py"
|
||||
return False
|
||||
|
||||
iso = IsolatedPlugin(name, str(plugin_path), {"settings": {"isolation": True}})
|
||||
ps.plugins[name] = iso
|
||||
logger.info(f"🔌 插件已导入(隔离模式): {name}")
|
||||
return True
|
||||
except Exception as e:
|
||||
_IMPORT_FAILS[name] = str(e)[:200]
|
||||
logger.error(f"插件加载失败 {name}: {e}")
|
||||
return False
|
||||
|
||||
|
||||
# ── 列表 / 管理 / 权限 ──
|
||||
|
||||
async def list_plugins(req):
|
||||
sm = req.app.get('service_manager')
|
||||
sm = req.app.get("service_manager")
|
||||
if not sm:
|
||||
return web.json_response({"error": "Service Manager 未初始化"}, status=503)
|
||||
|
||||
ps = sm.get_service("plugin")
|
||||
if not ps:
|
||||
return web.json_response({"plugins": []})
|
||||
|
||||
data = []
|
||||
for name, info in ps.plugin_info.items():
|
||||
data.append({
|
||||
"name": name,
|
||||
"version": getattr(info, 'version', '?'),
|
||||
"version": getattr(info, "version", "?"),
|
||||
"running": name in ps.plugins,
|
||||
"enabled": True
|
||||
"enabled": getattr(info, "enabled", True),
|
||||
"error": _IMPORT_FAILS.get(name, ""),
|
||||
})
|
||||
return web.json_response({"plugins": data})
|
||||
|
||||
|
||||
async def manage_plugin(req):
|
||||
sm = req.app.get('service_manager')
|
||||
if not sm: return web.json_response({"error": "SM Missing"}, 503)
|
||||
|
||||
name = req.match_info['name']
|
||||
action = req.match_info['action']
|
||||
sm = req.app.get("service_manager")
|
||||
if not sm:
|
||||
return web.json_response({"error": "SM Missing"}, status=503)
|
||||
name = req.match_info["name"]
|
||||
action = req.match_info["action"]
|
||||
ps = sm.get_service("plugin")
|
||||
|
||||
if not ps: return web.json_response({"error": "Plugin Service Missing"}, 503)
|
||||
|
||||
if not ps:
|
||||
return web.json_response({"error": "Plugin Service Missing"}, status=503)
|
||||
try:
|
||||
if action in ('disable', 'unload'):
|
||||
if action in ("disable", "unload"):
|
||||
await ps.unload_plugin(name)
|
||||
elif action == 'enable':
|
||||
elif action == "enable":
|
||||
await ps.load_plugin(name)
|
||||
elif action == 'reload':
|
||||
elif action == "reload":
|
||||
await ps.unload_plugin(name)
|
||||
await ps.load_plugin(name)
|
||||
return web.json_response({"success": True, "msg": "操作成功"})
|
||||
@@ -55,9 +246,9 @@ async def manage_plugin(req):
|
||||
logger.error(f"插件操作失败: {e}")
|
||||
return web.json_response({"success": False, "error": str(e)})
|
||||
|
||||
|
||||
async def list_plugin_web_pages(req):
|
||||
"""Return all registered plugin web UI pages for sidebar listing."""
|
||||
sm = req.app.get('service_manager')
|
||||
sm = req.app.get("service_manager")
|
||||
if not sm:
|
||||
return web.json_response({"pages": []})
|
||||
ps = sm.get_service("plugin")
|
||||
@@ -75,8 +266,10 @@ async def list_plugin_web_pages(req):
|
||||
})
|
||||
return web.json_response({"pages": pages})
|
||||
|
||||
|
||||
async def get_perms(req):
|
||||
return web.json_response({"plugin": req.match_info['name'], "permissions": ["read", "write"]})
|
||||
return web.json_response({"plugin": req.match_info["name"], "permissions": ["read", "write"]})
|
||||
|
||||
|
||||
async def set_perms(req):
|
||||
return web.json_response({"success": True})
|
||||
|
||||
Reference in New Issue
Block a user