Coverage for src\sync_monitor.py: 0%

184 statements  

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

1""" 

2Sync Monitor Module - 同步状态监控模块 

3 

4提供同步状态跟踪、性能指标、日志记录等功能 

5""" 

6 

7import asyncio 

8import logging 

9import time 

10from typing import Dict, List, Optional, Any 

11from dataclasses import dataclass, field 

12from enum import Enum 

13from collections import deque 

14import json 

15 

16logger = logging.getLogger(__name__) 

17 

18 

19class SyncState(Enum): 

20 """同步状态""" 

21 IDLE = "idle" 

22 SYNCING = "syncing" 

23 SYNCED = "synced" 

24 ERROR = "error" 

25 OFFLINE = "offline" 

26 

27 

28class SyncEventType(Enum): 

29 """同步事件类型""" 

30 START = "sync_start" 

31 COMPLETE = "sync_complete" 

32 ERROR = "sync_error" 

33 CONFLICT = "sync_conflict" 

34 RETRY = "sync_retry" 

35 

36 

37@dataclass 

38class SyncEvent: 

39 """同步事件""" 

40 event_type: SyncEventType 

41 timestamp: float 

42 peer_id: str 

43 details: Dict[str, Any] = field(default_factory=dict) 

44 duration_ms: float = 0 

45 

46 

47@dataclass 

48class SyncMetrics: 

49 """同步性能指标""" 

50 total_syncs: int = 0 

51 successful_syncs: int = 0 

52 failed_syncs: int = 0 

53 conflicts_resolved: int = 0 

54 total_bytes_sent: int = 0 

55 total_bytes_received: int = 0 

56 avg_latency_ms: float = 0 

57 last_sync_time: float = 0 

58 last_sync_duration_ms: float = 0 

59 

60 

61class SyncStateTracker: 

62 """同步状态跟踪器""" 

63 

64 def __init__(self): 

65 self._state: SyncState = SyncState.IDLE 

66 self._peer_states: Dict[str, SyncState] = {} 

67 self._state_history: deque = deque(maxlen=100) 

68 self._callbacks: Dict[str, List[callable]] = {} 

69 self._lock = asyncio.Lock() 

70 

71 @property 

72 def state(self) -> SyncState: 

73 """获取当前状态""" 

74 return self._state 

75 

76 async def set_state(self, state: SyncState, peer_id: Optional[str] = None) -> None: 

77 """设置同步状态""" 

78 async with self._lock: 

79 old_state = self._state 

80 self._state = state 

81 

82 # 记录历史 

83 self._state_history.append({ 

84 "from": old_state.value, 

85 "to": state.value, 

86 "peer_id": peer_id, 

87 "timestamp": time.time() 

88 }) 

89 

90 # 触发回调 

91 await self._trigger_callback("state_changed", old_state, state, peer_id) 

92 

93 logger.info(f"同步状态变化: {old_state.value} -> {state.value}" + 

94 (f" (peer: {peer_id})" if peer_id else "")) 

95 

96 async def set_peer_state(self, peer_id: str, state: SyncState) -> None: 

97 """设置对等节点状态""" 

98 async with self._lock: 

99 self._peer_states[peer_id] = state 

100 await self._trigger_callback("peer_state_changed", peer_id, state) 

101 

102 async def get_peer_state(self, peer_id: str) -> SyncState: 

103 """获取对等节点状态""" 

104 return self._peer_states.get(peer_id, SyncState.OFFLINE) 

105 

106 async def get_all_peer_states(self) -> Dict[str, SyncState]: 

107 """获取所有对等节点状态""" 

108 return self._peer_states.copy() 

109 

110 def get_state_history(self) -> List[Dict]: 

111 """获取状态历史""" 

112 return list(self._state_history) 

113 

114 def register_callback(self, event: str, callback: callable) -> None: 

115 """注册回调""" 

116 if event not in self._callbacks: 

117 self._callbacks[event] = [] 

118 self._callbacks[event].append(callback) 

119 

120 async def _trigger_callback(self, event: str, *args, **kwargs) -> None: 

121 """触发回调""" 

122 if event in self._callbacks: 

123 for callback in self._callbacks[event]: 

124 try: 

125 if asyncio.iscoroutinefunction(callback): 

126 await callback(*args, **kwargs) 

127 else: 

128 callback(*args, **kwargs) 

129 except Exception as e: 

130 logger.error(f"回调执行失败: {e}") 

131 

132 

133class SyncMetricsCollector: 

134 """同步性能指标收集器""" 

135 

136 def __init__(self): 

137 self._metrics = SyncMetrics() 

138 self._latencies: deque = deque(maxlen=1000) 

139 self._lock = asyncio.Lock() 

140 

141 async def record_sync_start(self) -> None: 

142 """记录同步开始""" 

143 async with self._lock: 

144 self._metrics.total_syncs += 1 

145 

146 async def record_sync_complete(self, duration_ms: float, bytes_sent: int = 0, bytes_received: int = 0) -> None: 

147 """记录同步完成""" 

148 async with self._lock: 

149 self._metrics.successful_syncs += 1 

150 self._metrics.last_sync_time = time.time() 

151 self._metrics.last_sync_duration_ms = duration_ms 

152 self._metrics.total_bytes_sent += bytes_sent 

153 self._metrics.total_bytes_received += bytes_received 

154 

155 # 更新平均延迟 

156 self._latencies.append(duration_ms) 

157 self._metrics.avg_latency_ms = sum(self._latencies) / len(self._latencies) 

158 

159 async def record_sync_error(self) -> None: 

160 """记录同步错误""" 

161 async with self._lock: 

162 self._metrics.failed_syncs += 1 

163 

164 async def record_conflict(self) -> None: 

165 """记录冲突""" 

166 async with self._lock: 

167 self._metrics.conflicts_resolved += 1 

168 

169 async def get_metrics(self) -> Dict[str, Any]: 

170 """获取指标""" 

171 async with self._lock: 

172 return { 

173 "total_syncs": self._metrics.total_syncs, 

174 "successful_syncs": self._metrics.successful_syncs, 

175 "failed_syncs": self._metrics.failed_syncs, 

176 "conflicts_resolved": self._metrics.conflicts_resolved, 

177 "total_bytes_sent": self._metrics.total_bytes_sent, 

178 "total_bytes_received": self._metrics.total_bytes_received, 

179 "avg_latency_ms": round(self._metrics.avg_latency_ms, 2), 

180 "last_sync_time": self._metrics.last_sync_time, 

181 "last_sync_duration_ms": round(self._metrics.last_sync_duration_ms, 2), 

182 "success_rate": self._calculate_success_rate() 

183 } 

184 

185 def _calculate_success_rate(self) -> float: 

186 """计算成功率""" 

187 if self._metrics.total_syncs == 0: 

188 return 0.0 

189 return self._metrics.successful_syncs / self._metrics.total_syncs * 100 

190 

191 async def reset(self) -> None: 

192 """重置指标""" 

193 async with self._lock: 

194 self._metrics = SyncMetrics() 

195 self._latencies.clear() 

196 

197 

198class SyncLogger: 

199 """同步日志记录器""" 

200 

201 def __init__(self, log_file: Optional[str] = None): 

202 self._log_file = log_file 

203 self._events: deque = deque(maxlen=1000) 

204 self._lock = asyncio.Lock() 

205 

206 async def log_event(self, event: SyncEvent) -> None: 

207 """记录事件""" 

208 async with self._lock: 

209 self._events.append(event) 

210 

211 # 同时输出到标准日志 

212 log_msg = f"[SYNC] {event.event_type.value} - peer: {event.peer_id}" 

213 if event.details: 

214 log_msg += f" - {json.dumps(event.details)}" 

215 if event.duration_ms > 0: 

216 log_msg += f" - {event.duration_ms:.2f}ms" 

217 

218 logger.info(log_msg) 

219 

220 # 写入文件 

221 if self._log_file: 

222 await self._write_to_file(event) 

223 

224 async def _write_to_file(self, event: SyncEvent) -> None: 

225 """写入日志文件""" 

226 try: 

227 with open(self._log_file, 'a', encoding='utf-8') as f: 

228 log_entry = { 

229 "event_type": event.event_type.value, 

230 "timestamp": event.timestamp, 

231 "peer_id": event.peer_id, 

232 "details": event.details, 

233 "duration_ms": event.duration_ms 

234 } 

235 f.write(json.dumps(log_entry, ensure_ascii=False) + "\n") 

236 except Exception as e: 

237 logger.error(f"写入日志文件失败: {e}") 

238 

239 async def get_recent_events(self, count: int = 50) -> List[SyncEvent]: 

240 """获取最近事件""" 

241 async with self._lock: 

242 return list(self._events)[-count:] 

243 

244 async def get_events_by_peer(self, peer_id: str) -> List[SyncEvent]: 

245 """获取指定对等节点的事件""" 

246 async with self._lock: 

247 return [e for e in self._events if e.peer_id == peer_id] 

248 

249 async def get_events_by_type(self, event_type: SyncEventType) -> List[SyncEvent]: 

250 """获取指定类型的事件""" 

251 async with self._lock: 

252 return [e for e in self._events if e.event_type == event_type] 

253 

254 async def clear(self) -> None: 

255 """清空日志""" 

256 async with self._lock: 

257 self._events.clear() 

258 

259 

260class SyncMonitor: 

261 """同步监控器(整合模块)""" 

262 

263 def __init__(self, log_file: Optional[str] = None): 

264 self.state_tracker = SyncStateTracker() 

265 self.metrics_collector = SyncMetricsCollector() 

266 self.logger = SyncLogger(log_file) 

267 self._running = False 

268 

269 async def start(self) -> None: 

270 """启动监控""" 

271 self._running = True 

272 logger.info("同步监控器已启动") 

273 

274 async def stop(self) -> None: 

275 """停止监控""" 

276 self._running = False 

277 logger.info("同步监控器已停止") 

278 

279 async def on_sync_start(self, peer_id: str) -> None: 

280 """同步开始回调""" 

281 await self.state_tracker.set_state(SyncState.SYNCING, peer_id) 

282 await self.state_tracker.set_peer_state(peer_id, SyncState.SYNCING) 

283 await self.metrics_collector.record_sync_start() 

284 

285 event = SyncEvent( 

286 event_type=SyncEventType.START, 

287 timestamp=time.time(), 

288 peer_id=peer_id 

289 ) 

290 await self.logger.log_event(event) 

291 

292 async def on_sync_complete(self, peer_id: str, duration_ms: float, bytes_sent: int = 0, bytes_received: int = 0) -> None: 

293 """同步完成回调""" 

294 await self.state_tracker.set_state(SyncState.SYNCED, peer_id) 

295 await self.state_tracker.set_peer_state(peer_id, SyncState.SYNCED) 

296 await self.metrics_collector.record_sync_complete(duration_ms, bytes_sent, bytes_received) 

297 

298 event = SyncEvent( 

299 event_type=SyncEventType.COMPLETE, 

300 timestamp=time.time(), 

301 peer_id=peer_id, 

302 duration_ms=duration_ms 

303 ) 

304 await self.logger.log_event(event) 

305 

306 async def on_sync_error(self, peer_id: str, error: str) -> None: 

307 """同步错误回调""" 

308 await self.state_tracker.set_state(SyncState.ERROR, peer_id) 

309 await self.state_tracker.set_peer_state(peer_id, SyncState.ERROR) 

310 await self.metrics_collector.record_sync_error() 

311 

312 event = SyncEvent( 

313 event_type=SyncEventType.ERROR, 

314 timestamp=time.time(), 

315 peer_id=peer_id, 

316 details={"error": error} 

317 ) 

318 await self.logger.log_event(event) 

319 

320 async def on_conflict(self, peer_id: str, conflict_type: str, resolution: str) -> None: 

321 """冲突解决回调""" 

322 await self.metrics_collector.record_conflict() 

323 

324 event = SyncEvent( 

325 event_type=SyncEventType.CONFLICT, 

326 timestamp=time.time(), 

327 peer_id=peer_id, 

328 details={"conflict_type": conflict_type, "resolution": resolution} 

329 ) 

330 await self.logger.log_event(event) 

331 

332 async def get_full_status(self) -> Dict[str, Any]: 

333 """获取完整状态""" 

334 return { 

335 "state": self.state_tracker.state.value, 

336 "peer_states": { 

337 peer_id: state.value 

338 for peer_id, state in (await self.state_tracker.get_all_peer_states()).items() 

339 }, 

340 "metrics": await self.metrics_collector.get_metrics(), 

341 "running": self._running 

342 }