8dbf628dd6
- 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)
154 lines
5.8 KiB
Python
154 lines
5.8 KiB
Python
#!/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 已关闭")
|