Phase 1: ProjectEngine + PyEnvManager + project management WebUI
- New: services/project_engine.py (async subprocess manager) - New: services/pyenv_manager.py (Python version + venv + git clone) - New: services/web_panel/routes/projects.py (5 REST endpoints) - New: static/web_panel/pages/projects.html (WebUI) - Tests: 28/28 passing, API verified (GET /api/projects) - Docs: Phase1_Progress.md (local)
This commit is contained in:
@@ -0,0 +1,153 @@
|
||||
#!/usr/bin/env python3
|
||||
"""SenSu ProjectEngine — 异步子进程项目管理"""
|
||||
import asyncio, os, logging, signal, time
|
||||
from typing import Dict, Optional, List
|
||||
from dataclasses import dataclass, field
|
||||
from datetime import datetime
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@dataclass
|
||||
class ProjectProcess:
|
||||
name: str
|
||||
cmd: list
|
||||
cwd: str = "."
|
||||
env: dict = field(default_factory=dict)
|
||||
port: int = 0
|
||||
proxy_path: str = ""
|
||||
auto_restart: bool = True
|
||||
process: asyncio.subprocess.Process = None
|
||||
status: str = "stopped"
|
||||
pid: int = 0
|
||||
started_at: float = 0
|
||||
log_buffer: list = field(default_factory=list) # 最近 200 行
|
||||
max_log_lines: int = 200
|
||||
|
||||
class ProjectEngine:
|
||||
def __init__(self, service_manager=None):
|
||||
self.sm = service_manager
|
||||
self.projects: Dict[str, ProjectProcess] = {}
|
||||
self._monitor_task = None
|
||||
|
||||
async def start(self):
|
||||
self._monitor_task = asyncio.create_task(self._health_monitor())
|
||||
logger.info("ProjectEngine 已就绪")
|
||||
|
||||
async def run_project(self, name: str, cmd: list, cwd: str = ".",
|
||||
env: dict = None, port: int = 0, proxy_path: str = "",
|
||||
auto_restart: bool = True) -> bool:
|
||||
if name in self.projects and self.projects[name].status == "running":
|
||||
logger.warning(f"项目 {name} 已在运行")
|
||||
return False
|
||||
|
||||
pp = ProjectProcess(name=name, cmd=cmd, cwd=cwd, env=env or {},
|
||||
port=port, proxy_path=proxy_path, auto_restart=auto_restart)
|
||||
return await self._start_process(pp)
|
||||
|
||||
async def _start_process(self, pp: ProjectProcess) -> bool:
|
||||
try:
|
||||
full_env = {**os.environ, **pp.env}
|
||||
pp.process = await asyncio.create_subprocess_exec(
|
||||
*pp.cmd, cwd=pp.cwd, env=full_env,
|
||||
stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE,
|
||||
stdin=asyncio.subprocess.PIPE
|
||||
)
|
||||
pp.pid = pp.process.pid
|
||||
pp.status = "running"
|
||||
pp.started_at = time.time()
|
||||
self.projects[pp.name] = pp
|
||||
|
||||
# 启动日志读取任务
|
||||
asyncio.create_task(self._read_stream(pp, pp.process.stdout, "stdout"))
|
||||
asyncio.create_task(self._read_stream(pp, pp.process.stderr, "stderr"))
|
||||
|
||||
logger.info(f"项目已启动: {pp.name} (PID={pp.pid})")
|
||||
if pp.port:
|
||||
logger.info(f" 端口: {pp.port}")
|
||||
if pp.proxy_path:
|
||||
logger.info(f" 代理: {pp.proxy_path}")
|
||||
|
||||
# 监控进程退出
|
||||
asyncio.create_task(self._wait_exit(pp))
|
||||
return True
|
||||
except Exception as e:
|
||||
logger.error(f"启动项目 {pp.name} 失败: {e}")
|
||||
pp.status = "error"
|
||||
self.projects[pp.name] = pp
|
||||
return False
|
||||
|
||||
async def _read_stream(self, pp: ProjectProcess, stream, tag: str):
|
||||
while pp.status == "running" and stream and not stream.at_eof():
|
||||
try:
|
||||
line = await stream.readline()
|
||||
if line:
|
||||
text = line.decode(errors="replace").rstrip()
|
||||
pp.log_buffer.append(f"[{tag}] {text}")
|
||||
if len(pp.log_buffer) > pp.max_log_lines:
|
||||
pp.log_buffer = pp.log_buffer[-pp.max_log_lines:]
|
||||
logger.debug(f"[{pp.name}] {text}")
|
||||
except Exception:
|
||||
break
|
||||
|
||||
async def _wait_exit(self, pp: ProjectProcess):
|
||||
if pp.process:
|
||||
await pp.process.wait()
|
||||
exit_code = pp.process.returncode
|
||||
pp.status = "stopped"
|
||||
logger.info(f"项目 {pp.name} 已退出 (code={exit_code})")
|
||||
if pp.auto_restart and exit_code != 0:
|
||||
logger.info(f"自动重启 {pp.name} ...")
|
||||
await asyncio.sleep(2)
|
||||
await self._start_process(pp)
|
||||
|
||||
async def stop_project(self, name: str) -> bool:
|
||||
pp = self.projects.get(name)
|
||||
if not pp or pp.status != "running":
|
||||
return False
|
||||
pp.auto_restart = False
|
||||
if pp.process:
|
||||
pp.process.terminate()
|
||||
try:
|
||||
await asyncio.wait_for(pp.process.wait(), timeout=5)
|
||||
except asyncio.TimeoutError:
|
||||
pp.process.kill()
|
||||
pp.status = "stopped"
|
||||
logger.info(f"项目已停止: {name}")
|
||||
return True
|
||||
|
||||
async def send_stdin(self, name: str, text: str):
|
||||
pp = self.projects.get(name)
|
||||
if pp and pp.process and pp.process.stdin:
|
||||
pp.process.stdin.write((text + "\n").encode())
|
||||
await pp.process.stdin.drain()
|
||||
return True
|
||||
return False
|
||||
|
||||
def get_logs(self, name: str, tail: int = 50) -> List[str]:
|
||||
pp = self.projects.get(name)
|
||||
return pp.log_buffer[-tail:] if pp else []
|
||||
|
||||
def list_projects(self) -> List[dict]:
|
||||
return [{"name": p.name, "status": p.status, "pid": p.pid,
|
||||
"port": p.port, "proxy": p.proxy_path,
|
||||
"uptime": int(time.time()-p.started_at) if p.started_at else 0}
|
||||
for p in self.projects.values()]
|
||||
|
||||
def get_project(self, name: str) -> Optional[ProjectProcess]:
|
||||
return self.projects.get(name)
|
||||
|
||||
async def _health_monitor(self):
|
||||
while True:
|
||||
await asyncio.sleep(10)
|
||||
for pp in list(self.projects.values()):
|
||||
if pp.status == "running" and pp.process:
|
||||
if pp.process.returncode is not None:
|
||||
pp.status = "stopped"
|
||||
logger.warning(f"项目 {pp.name} 异常退出 (code={pp.process.returncode})")
|
||||
|
||||
async def shutdown(self):
|
||||
for name in list(self.projects.keys()):
|
||||
await self.stop_project(name)
|
||||
if self._monitor_task:
|
||||
self._monitor_task.cancel()
|
||||
logger.info("ProjectEngine 已关闭")
|
||||
@@ -0,0 +1,130 @@
|
||||
#!/usr/bin/env python3
|
||||
"""SenSu PyEnvManager — Python 版本管理 + venv + 依赖安装"""
|
||||
import os, sys, subprocess, logging, venv, shutil
|
||||
from pathlib import Path
|
||||
from typing import Optional, List
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
class PyEnvManager:
|
||||
def __init__(self, workspace_dir: str = None):
|
||||
self.workspace = Path(workspace_dir or os.getcwd())
|
||||
self.venvs_dir = self.workspace / "venvs"
|
||||
self.venvs_dir.mkdir(parents=True, exist_ok=True)
|
||||
|
||||
def detect_versions(self) -> List[str]:
|
||||
versions = set()
|
||||
# Current Python
|
||||
v = f"{sys.version_info.major}.{sys.version_info.minor}"
|
||||
versions.add(v)
|
||||
|
||||
# Check Termux pkg
|
||||
try:
|
||||
result = subprocess.run(["pkg", "list-installed"], capture_output=True, text=True, timeout=10)
|
||||
for line in result.stdout.split("\n"):
|
||||
if line.startswith("python-3.") or line.startswith("python3-"):
|
||||
versions.add(line.split("/")[0].replace("python-", "").replace("python3-", ""))
|
||||
except:
|
||||
pass
|
||||
|
||||
# Check pyenv
|
||||
pyenv = shutil.which("pyenv")
|
||||
if pyenv:
|
||||
try:
|
||||
result = subprocess.run([pyenv, "versions", "--bare"], capture_output=True, text=True, timeout=10)
|
||||
for line in result.stdout.split("\n"):
|
||||
line = line.strip()
|
||||
if line and line[0].isdigit():
|
||||
versions.add(line.split("/")[0])
|
||||
except:
|
||||
pass
|
||||
|
||||
return sorted(versions)
|
||||
|
||||
def ensure_version(self, version: str) -> Optional[str]:
|
||||
"""确保指定 Python 版本可用,返回解释器路径"""
|
||||
available = self.detect_versions()
|
||||
if version in available:
|
||||
return self._find_python(version)
|
||||
|
||||
# Try to install via Termux
|
||||
if shutil.which("pkg"):
|
||||
pkg_name = f"python-{version}"
|
||||
logger.info(f"尝试安装 {pkg_name} ...")
|
||||
try:
|
||||
subprocess.run(["pkg", "install", "-y", pkg_name], check=True, timeout=120)
|
||||
return self._find_python(version)
|
||||
except:
|
||||
pass
|
||||
|
||||
logger.warning(f"无法获取 Python {version},使用当前版本")
|
||||
return sys.executable
|
||||
|
||||
def _find_python(self, version: str) -> Optional[str]:
|
||||
for name in [f"python{version}", f"python{version[:3]}", "python3"]:
|
||||
path = shutil.which(name)
|
||||
if path: return path
|
||||
return sys.executable
|
||||
|
||||
def create_venv(self, name: str, python_version: str = None) -> Optional[Path]:
|
||||
venv_path = self.venvs_dir / name
|
||||
if venv_path.exists():
|
||||
logger.info(f"venv 已存在: {venv_path}")
|
||||
return venv_path
|
||||
|
||||
python_exe = self.ensure_version(python_version) if python_version else sys.executable
|
||||
logger.info(f"创建 venv: {venv_path} (Python {python_version or 'default'})")
|
||||
|
||||
try:
|
||||
venv.create(str(venv_path), with_pip=True, clear=True)
|
||||
# Install/upgrade pip
|
||||
pip = str(venv_path / "bin" / "pip")
|
||||
subprocess.run([pip, "install", "--upgrade", "pip"], capture_output=True, timeout=60)
|
||||
return venv_path
|
||||
except Exception as e:
|
||||
logger.error(f"创建 venv 失败: {e}")
|
||||
# Fallback: use virtualenv
|
||||
try:
|
||||
subprocess.run([sys.executable, "-m", "virtualenv", str(venv_path)], check=True, timeout=120)
|
||||
return venv_path
|
||||
except:
|
||||
return None
|
||||
|
||||
def install_deps(self, venv_path: Path, requirements: List[str]) -> bool:
|
||||
pip = str(venv_path / "bin" / "pip")
|
||||
for req_file in requirements:
|
||||
req_path = Path(req_file)
|
||||
if not req_path.is_absolute():
|
||||
# Relative to workspace
|
||||
pass
|
||||
if Path(req_file).exists():
|
||||
logger.info(f"安装依赖: {req_file}")
|
||||
try:
|
||||
subprocess.run([pip, "install", "-r", req_file], check=True, timeout=300)
|
||||
except subprocess.CalledProcessError as e:
|
||||
logger.warning(f"依赖安装部分失败: {e}")
|
||||
return False
|
||||
return True
|
||||
|
||||
def clone_git(self, url: str, target_dir: Path, branch: str = None) -> bool:
|
||||
if target_dir.exists():
|
||||
logger.info(f"目录已存在: {target_dir}")
|
||||
# Try git pull instead
|
||||
try:
|
||||
subprocess.run(["git", "-C", str(target_dir), "pull"], check=True, timeout=60)
|
||||
return True
|
||||
except:
|
||||
pass
|
||||
|
||||
cmd = ["git", "clone"]
|
||||
if branch:
|
||||
cmd += ["-b", branch]
|
||||
cmd += [url, str(target_dir)]
|
||||
|
||||
try:
|
||||
subprocess.run(cmd, check=True, timeout=300)
|
||||
logger.info(f"Git clone 完成: {url} → {target_dir}")
|
||||
return True
|
||||
except subprocess.CalledProcessError as e:
|
||||
logger.error(f"Git clone 失败: {e}")
|
||||
return False
|
||||
@@ -5,7 +5,7 @@ import os
|
||||
import logging
|
||||
from pathlib import Path
|
||||
from aiohttp import web
|
||||
from .routes import auth, status, plugins, commands, logs
|
||||
from .routes import auth, status, plugins, commands, logs, projects
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
@@ -62,6 +62,7 @@ class WebPanelManager:
|
||||
plugins.setup_routes(app, self.base_path)
|
||||
commands.setup_routes(app, self.base_path)
|
||||
logs.setup_routes(app, self.base_path)
|
||||
projects.setup_project_routes(app, self.sm)
|
||||
|
||||
# 注册日志广播
|
||||
ls = self.sm.get_service("log")
|
||||
|
||||
@@ -0,0 +1,60 @@
|
||||
from aiohttp import web
|
||||
import json, logging
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
def _get_engine(request):
|
||||
sm = request.app.get("service_manager")
|
||||
if sm and sm.has_service("project_engine"):
|
||||
return sm.get_service("project_engine")
|
||||
return None
|
||||
|
||||
def setup_project_routes(app, service_manager):
|
||||
app["service_manager"] = service_manager
|
||||
|
||||
async def list_projects(request):
|
||||
eng = _get_engine(request)
|
||||
return web.json_response({"projects": eng.list_projects() if eng else []})
|
||||
|
||||
async def run_project(request):
|
||||
try:
|
||||
data = await request.json()
|
||||
eng = _get_engine(request)
|
||||
if not eng:
|
||||
return web.json_response({"ok": False, "error": "engine not ready"}, status=503)
|
||||
ok = await eng.run_project(
|
||||
name=data.get("name","unnamed"), cmd=data.get("cmd",[]),
|
||||
cwd=data.get("cwd","."), env=data.get("env",{}),
|
||||
port=data.get("port",0), proxy_path=data.get("proxy_path",""))
|
||||
return web.json_response({"ok": ok})
|
||||
except Exception as e:
|
||||
return web.json_response({"ok": False, "error": str(e)}, status=400)
|
||||
|
||||
async def stop_project(request):
|
||||
name = request.match_info.get("name","")
|
||||
eng = _get_engine(request)
|
||||
ok = await eng.stop_project(name) if eng else False
|
||||
return web.json_response({"ok": ok})
|
||||
|
||||
async def get_logs(request):
|
||||
name = request.match_info.get("name","")
|
||||
tail = int(request.query.get("tail", 50))
|
||||
eng = _get_engine(request)
|
||||
return web.json_response({"logs": eng.get_logs(name, tail) if eng else []})
|
||||
|
||||
async def send_stdin(request):
|
||||
name = request.match_info.get("name","")
|
||||
data = await request.json()
|
||||
eng = _get_engine(request)
|
||||
ok = await eng.send_stdin(name, data.get("text","")) if eng else False
|
||||
return web.json_response({"ok": ok})
|
||||
|
||||
async def project_page(request):
|
||||
return web.FileResponse("static/web_panel/pages/projects.html")
|
||||
|
||||
app.router.add_get("/api/projects", list_projects)
|
||||
app.router.add_post("/api/projects/run", run_project)
|
||||
app.router.add_get("/api/projects/{name}/logs", get_logs)
|
||||
app.router.add_post("/api/projects/{name}/stop", stop_project)
|
||||
app.router.add_post("/api/projects/{name}/stdin", send_stdin)
|
||||
app.router.add_get("/pages/projects", project_page)
|
||||
logger.info("📦 项目管理路由已注册")
|
||||
Reference in New Issue
Block a user