d48ef65b35
- 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)
118 lines
4.3 KiB
Python
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
|