Files
SenSu/services/project_service.py
T
AskaEth d48ef65b35 v0.3: SQLite persistence + ProjectService + dependency resolution
- New: services/sensu_db.py (SQLite, 5 tables)
- New: services/project_service.py (project registry, port allocation, dep resolution)
- Integrated into main.py (step 11.6)
- Tests: 23/23 passing (15 original + 8 new)
2026-06-10 18:59:25 +08:00

118 lines
4.3 KiB
Python

#!/usr/bin/env python3
"""SenSu 项目注册表 — 插件声明 project.yaml 申请资源,框架管理生命周期"""
import logging, os, yaml, asyncio
from typing import Dict, List, Optional
from collections import defaultdict, deque
logger = logging.getLogger(__name__)
class ProjectService:
def __init__(self, service_manager, db=None):
self.sm = service_manager
self.db = db
self.projects: Dict[str, dict] = {}
self._port_allocations: Dict[int, str] = {}
async def start(self):
logger.info("项目注册表已就绪")
def register_project(self, plugin_name: str, project_config: dict) -> bool:
"""从 project.yaml 注册项目
必需字段: name, port (可选: entrypoint, depends_on, env)
"""
name = project_config.get("name", plugin_name)
if name in self.projects:
logger.warning(f"项目 {name} 已注册,跳过")
return False
port = project_config.get("port")
if port and port in self._port_allocations:
logger.warning(f"端口 {port} 已被 {self._port_allocations[port]} 占用")
port = self._find_free_port()
self.projects[name] = {
**project_config,
"plugin_name": plugin_name,
"port": port,
"status": "registered",
}
if port:
self._port_allocations[port] = name
if self.db:
self.db.save_plugin(plugin_name, project_path=project_config.get("path",""),
project_port=port or 0)
self.db.log_audit(plugin_name, "project_registered", str(project_config))
logger.info(f"项目已注册: {name} (插件: {plugin_name}, 端口: {port})")
return True
def _find_free_port(self, start=4200) -> int:
used = set(self._port_allocations.keys())
for p in range(start, start + 1000):
if p not in used:
return p
return start
def get_project(self, name: str) -> Optional[dict]:
return self.projects.get(name)
def list_projects(self) -> List[dict]:
return [{"name": k, "port": v.get("port"), "status": v.get("status")}
for k, v in self.projects.items()]
def unregister_project(self, name: str):
if name in self.projects:
p = self.projects.pop(name)
if p.get("port"):
self._port_allocations.pop(p["port"], None)
logger.info(f"项目已注销: {name}")
def resolve_dependencies(self, plugins: Dict[str, dict]) -> List[str]:
"""拓扑排序 — 根据 depends_on 返回正确的加载顺序"""
graph = defaultdict(list)
in_degree = defaultdict(int)
all_plugins = set(plugins.keys())
for name, cfg in plugins.items():
deps = cfg.get("depends_on", [])
if isinstance(deps, str):
deps = json.loads(deps) if deps.startswith("[") else [deps]
for dep in deps:
if dep in all_plugins:
graph[dep].append(name)
in_degree[name] += 1
if name not in in_degree:
in_degree[name] = 0
# Kahn's algorithm
queue = deque([n for n in all_plugins if in_degree[n] == 0])
result = []
while queue:
node = queue.popleft()
result.append(node)
for neighbor in graph[node]:
in_degree[neighbor] -= 1
if in_degree[neighbor] == 0:
queue.append(neighbor)
if len(result) != len(all_plugins):
missing = all_plugins - set(result)
logger.warning(f"循环依赖或缺失依赖: {missing}, 追加到尾部")
result.extend(missing)
logger.info(f"依赖解析结果: {''.join(result)}")
return result
@staticmethod
def load_project_yaml(plugin_dir: str) -> Optional[dict]:
"""从插件目录加载 project.yaml"""
yaml_path = os.path.join(plugin_dir, "project.yaml")
if os.path.exists(yaml_path):
try:
with open(yaml_path) as f:
return yaml.safe_load(f)
except Exception as e:
logger.warning(f"解析 project.yaml 失败 {yaml_path}: {e}")
return None