Coverage for src\system_integrator.py: 0%

181 statements  

« prev     ^ index     » next       coverage.py v7.3.4, created at 2026-04-21 14:55 +0800

1""" 

2PAO System Integrator - 系统集成器 

3 

4整合所有模块,确保协同工作: 

5- 设备发现与通信 

6- 技能管理 

7- 情境感知 

8- 学习循环 

9- 数据同步 

10- 记忆系统 

11""" 

12 

13import asyncio 

14import logging 

15from typing import Optional, Dict, Any 

16from dataclasses import dataclass 

17 

18from .core.config import ConfigManager, PAOConfig 

19from .core.device import DeviceRegistry, DeviceInfo 

20from .core.discovery import DeviceDiscoveryService 

21from .core.communication import Communication 

22from .core.sync import SyncEngine, SyncStatus 

23from .core.storage import StorageManager, StorageConfig as CoreStorageConfig 

24from .core.memory import MemorySystem 

25from .core.heartbeat import HeartbeatManager, HeartbeatConfig 

26 

27from .skill_manager import SkillManager, SkillCategory, SkillLevel 

28from .context_awareness import ContextAwareness, ContextType 

29from .learning_loop import LearningLoop, FeedbackType 

30 

31logger = logging.getLogger(__name__) 

32 

33 

34@dataclass 

35class SystemStatus: 

36 """系统状态""" 

37 initialized: bool = False 

38 device_id: Optional[str] = None 

39 connected_peers: int = 0 

40 skills_loaded: int = 0 

41 memory_items: int = 0 

42 sync_status: str = "idle" 

43 errors: list = None 

44 

45 def __post_init__(self): 

46 if self.errors is None: 

47 self.errors = [] 

48 

49 

50class PAOSystemIntegrator: 

51 """ 

52 PAO系统集成器 

53 

54 整合所有核心模块,提供统一的系统入口 

55 """ 

56 

57 def __init__(self, config_path: Optional[str] = None): 

58 self.config_manager = ConfigManager(config_path) 

59 self.config: PAOConfig = self.config_manager.config 

60 

61 # 核心组件 

62 self.device_registry: Optional[DeviceRegistry] = None 

63 self.discovery_service: Optional[DeviceDiscoveryService] = None 

64 self.communication: Optional[Communication] = None 

65 self.sync_engine: Optional[SyncEngine] = None 

66 self.storage_manager: Optional[StorageManager] = None 

67 self.memory_system: Optional[MemorySystem] = None 

68 self.heartbeat_manager: Optional[HeartbeatManager] = None 

69 

70 # 智能组件 

71 self.skill_manager: Optional[SkillManager] = None 

72 self.context_awareness: Optional[ContextAwareness] = None 

73 self.learning_loop: Optional[LearningLoop] = None 

74 

75 # 系统状态 

76 self.status = SystemStatus() 

77 self._running = False 

78 

79 async def initialize(self) -> bool: 

80 """ 

81 初始化所有系统组件 

82 

83 Returns: 

84 bool: 初始化是否成功 

85 """ 

86 logger.info("🚀 初始化PAO系统...") 

87 

88 try: 

89 # 1. 初始化设备注册表 

90 self.device_registry = DeviceRegistry() 

91 self.status.device_id = self.config.device_id 

92 

93 # 2. 初始化存储管理器 

94 self.storage_manager = StorageManager(self.config.storage) 

95 await self.storage_manager.initialize() 

96 

97 # 3. 初始化记忆系统 

98 self.memory_system = MemorySystem() 

99 await self.memory_system.load() 

100 self.status.memory_items = len(self.memory_system.memories) 

101 

102 # 4. 初始化设备发现服务 

103 self.discovery_service = DeviceDiscoveryService( 

104 config=self.config, 

105 device_registry=self.device_registry, 

106 on_device_discovered=self._on_device_discovered, 

107 on_device_lost=self._on_device_lost 

108 ) 

109 

110 # 5. 初始化通信模块 

111 self.communication = Communication( 

112 device_id=self.config.device_id, 

113 config=self.config 

114 ) 

115 

116 # 6. 初始化同步引擎 

117 self.sync_engine = SyncEngine( 

118 memory_system=self.memory_system, 

119 local_device_id=self.config.device_id 

120 ) 

121 

122 # 7. 初始化技能管理器 

123 self.skill_manager = SkillManager() 

124 await self.skill_manager.initialize() 

125 self.status.skills_loaded = len(self.skill_manager.registry._skills) 

126 

127 # 8. 初始化情境感知 

128 self.context_awareness = ContextAwareness() 

129 await self.context_awareness.register_default_scenes() 

130 

131 # 9. 初始化学习循环 

132 self.learning_loop = LearningLoop() 

133 await self.learning_loop.start() 

134 

135 # 10. 初始化心跳管理器 

136 self.heartbeat_manager = HeartbeatManager( 

137 device_id=self.config.device_id, 

138 config=HeartbeatConfig( 

139 interval=30, # 30秒心跳间隔 

140 timeout=90, # 90秒超时判定离线 

141 check_interval=10 # 10秒检查间隔 

142 ) 

143 ) 

144 # 设置发送器(复用通信模块) 

145 self.heartbeat_manager.set_sender(self.communication) 

146 # 注册设备状态变化回调 

147 self.heartbeat_manager.register_callback(self._on_heartbeat_status_change) 

148 await self.heartbeat_manager.start() 

149 

150 # 启动设备发现 

151 await self.discovery_service.start() 

152 

153 self.status.initialized = True 

154 self._running = True 

155 

156 logger.info("✅ PAO系统初始化完成") 

157 logger.info(f" 设备ID: {self.status.device_id}") 

158 logger.info(f" 技能数量: {self.status.skills_loaded}") 

159 logger.info(f" 记忆数量: {self.status.memory_items}") 

160 

161 return True 

162 

163 except Exception as e: 

164 logger.error(f"❌ 系统初始化失败: {e}") 

165 self.status.errors.append(str(e)) 

166 return False 

167 

168 async def start(self): 

169 """启动系统""" 

170 if not self.status.initialized: 

171 success = await self.initialize() 

172 if not success: 

173 raise RuntimeError("系统初始化失败") 

174 

175 logger.info("▶️ 启动PAO系统服务...") 

176 

177 # 启动设备发现 

178 if self.discovery_service: 

179 await self.discovery_service.start() 

180 

181 # 启动同步引擎 

182 if self.sync_engine: 

183 await self.sync_engine.start() 

184 

185 self._running = True 

186 logger.info("✅ PAO系统已启动") 

187 

188 async def stop(self): 

189 """停止系统""" 

190 logger.info("⏹️ 停止PAO系统...") 

191 

192 self._running = False 

193 

194 # 停止各组件 

195 if self.heartbeat_manager: 

196 await self.heartbeat_manager.stop() 

197 

198 if self.discovery_service: 

199 await self.discovery_service.stop() 

200 

201 if self.sync_engine: 

202 await self.sync_engine.stop() 

203 

204 if self.learning_loop: 

205 await self.learning_loop.stop() 

206 

207 if self.storage_manager: 

208 await self.storage_manager.close() 

209 

210 logger.info("✅ PAO系统已停止") 

211 

212 # ==================== 事件回调 ==================== 

213 

214 def _on_device_discovered(self, device: DeviceInfo): 

215 """设备发现回调""" 

216 logger.info(f"📡 发现新设备: {device.name} ({device.device_id})") 

217 self.status.connected_peers = len(self.device_registry.list_devices()) 

218 

219 def _on_device_lost(self, device_id: str): 

220 """设备丢失回调""" 

221 logger.info(f"📡 设备丢失: {device_id}") 

222 self.status.connected_peers = len(self.device_registry.list_devices()) 

223 

224 def _on_sync_status_change(self, status: SyncStatus): 

225 """同步状态变化回调""" 

226 self.status.sync_status = status.value 

227 logger.debug(f"🔄 同步状态: {status.value}") 

228 

229 def _on_heartbeat_status_change(self, device_id: str, old_status, new_status): 

230 """心跳状态变化回调""" 

231 from .core.heartbeat import HeartbeatStatus 

232 

233 logger.info(f"💓 设备状态变化 [{device_id}]: {old_status.value} -> {new_status.value}") 

234 

235 # 如果设备离线,从设备注册表移除 

236 if new_status == HeartbeatStatus.DEAD: 

237 logger.warning(f"⚠️ 设备离线 [{device_id}],从注册表移除") 

238 self.device_registry.unregister_device(device_id) 

239 

240 # 更新连接数 

241 self.status.connected_peers = len(self.device_registry.list_online_devices()) 

242 

243 # ==================== 技能管理 ==================== 

244 

245 async def search_skills(self, query: str) -> list: 

246 """搜索技能""" 

247 if not self.skill_manager: 

248 return [] 

249 return await self.skill_manager.search_skills(query) 

250 

251 async def apply_skill(self, skill_id: str, params: Dict[str, Any], 

252 score: float, feedback: str) -> bool: 

253 """应用技能并评分""" 

254 if not self.skill_manager: 

255 return False 

256 

257 # 应用技能 

258 result = await self.skill_manager.apply_skill(skill_id, params, score, feedback) 

259 

260 # 提交学习反馈 

261 if self.learning_loop: 

262 await self.learning_loop.submit_feedback( 

263 FeedbackType.EXPLICIT, 

264 skill_id, 

265 params, 

266 score, 

267 feedback 

268 ) 

269 

270 return result 

271 

272 async def get_skill_stats(self, skill_id: str) -> Dict[str, Any]: 

273 """获取技能统计""" 

274 if not self.skill_manager: 

275 return {} 

276 return await self.skill_manager.get_skill_stats(skill_id) 

277 

278 # ==================== 情境感知 ==================== 

279 

280 async def get_current_context(self) -> Dict[str, Any]: 

281 """获取当前上下文""" 

282 if not self.context_awareness: 

283 return {} 

284 

285 contexts = await self.context_awareness.collect_all() 

286 scene = await self.context_awareness.recognize_scene() 

287 

288 return { 

289 "contexts": contexts, 

290 "scene": scene, 

291 "summary": await self.context_awareness.summarize_context() 

292 } 

293 

294 # ==================== 记忆管理 ==================== 

295 

296 async def store_memory(self, memory_type: str, content: Any, 

297 priority: int = 5) -> str: 

298 """存储记忆""" 

299 if not self.memory_system: 

300 return None 

301 

302 memory_id = await self.memory_system.add_memory( 

303 memory_type=memory_type, 

304 content=content, 

305 priority=priority 

306 ) 

307 self.status.memory_items = len(self.memory_system.memories) 

308 return memory_id 

309 

310 async def retrieve_memories(self, query: str, limit: int = 10) -> list: 

311 """检索记忆""" 

312 if not self.memory_system: 

313 return [] 

314 return await self.memory_system.search_memories(query, limit) 

315 

316 # ==================== 同步管理 ==================== 

317 

318 async def sync_with_peer(self, peer_id: str) -> bool: 

319 """与指定对等节点同步""" 

320 if not self.sync_engine: 

321 return False 

322 

323 try: 

324 await self.sync_engine.connect_peer(peer_id) 

325 await self.sync_engine.sync_now() 

326 return True 

327 except Exception as e: 

328 logger.error(f"同步失败: {e}") 

329 return False 

330 

331 def get_sync_status(self) -> Dict[str, Any]: 

332 """获取同步状态""" 

333 if not self.sync_engine: 

334 return {"status": "unavailable"} 

335 

336 status = self.sync_engine.get_status() 

337 return { 

338 "status": status.status.value if hasattr(status, 'status') else str(status), 

339 "connected_peers": status.connected_peers if hasattr(status, 'connected_peers') else 0, 

340 "pending_changes": status.pending_changes if hasattr(status, 'pending_changes') else 0 

341 } 

342 

343 # ==================== 系统状态 ==================== 

344 

345 def get_system_status(self) -> SystemStatus: 

346 """获取系统状态""" 

347 self.status.connected_peers = len(self.device_registry.get_all_devices()) if self.device_registry else 0 

348 return self.status 

349 

350 async def run_diagnostics(self) -> Dict[str, Any]: 

351 """运行系统诊断""" 

352 diagnostics = { 

353 "system": { 

354 "initialized": self.status.initialized, 

355 "running": self._running, 

356 "device_id": self.status.device_id 

357 }, 

358 "components": {}, 

359 "errors": self.status.errors 

360 } 

361 

362 # 检查各组件 

363 components = [ 

364 ("device_registry", self.device_registry), 

365 ("discovery_service", self.discovery_service), 

366 ("communication", self.communication), 

367 ("sync_engine", self.sync_engine), 

368 ("storage_manager", self.storage_manager), 

369 ("memory_system", self.memory_system), 

370 ("skill_manager", self.skill_manager), 

371 ("context_awareness", self.context_awareness), 

372 ("learning_loop", self.learning_loop) 

373 ] 

374 

375 for name, component in components: 

376 diagnostics["components"][name] = { 

377 "available": component is not None, 

378 "status": "ok" if component else "unavailable" 

379 } 

380 

381 return diagnostics 

382 

383 

384# ==================== 便捷函数 ==================== 

385 

386async def create_pao_system(config_path: Optional[str] = None) -> PAOSystemIntegrator: 

387 """创建并初始化PAO系统""" 

388 system = PAOSystemIntegrator(config_path) 

389 await system.initialize() 

390 return system