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
« prev ^ index » next coverage.py v7.3.4, created at 2026-04-16 16:31 +0800
1"""
2Sync Monitor Module - 同步状态监控模块
4提供同步状态跟踪、性能指标、日志记录等功能
5"""
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
16logger = logging.getLogger(__name__)
19class SyncState(Enum):
20 """同步状态"""
21 IDLE = "idle"
22 SYNCING = "syncing"
23 SYNCED = "synced"
24 ERROR = "error"
25 OFFLINE = "offline"
28class SyncEventType(Enum):
29 """同步事件类型"""
30 START = "sync_start"
31 COMPLETE = "sync_complete"
32 ERROR = "sync_error"
33 CONFLICT = "sync_conflict"
34 RETRY = "sync_retry"
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
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
61class SyncStateTracker:
62 """同步状态跟踪器"""
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()
71 @property
72 def state(self) -> SyncState:
73 """获取当前状态"""
74 return self._state
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
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 })
90 # 触发回调
91 await self._trigger_callback("state_changed", old_state, state, peer_id)
93 logger.info(f"同步状态变化: {old_state.value} -> {state.value}" +
94 (f" (peer: {peer_id})" if peer_id else ""))
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)
102 async def get_peer_state(self, peer_id: str) -> SyncState:
103 """获取对等节点状态"""
104 return self._peer_states.get(peer_id, SyncState.OFFLINE)
106 async def get_all_peer_states(self) -> Dict[str, SyncState]:
107 """获取所有对等节点状态"""
108 return self._peer_states.copy()
110 def get_state_history(self) -> List[Dict]:
111 """获取状态历史"""
112 return list(self._state_history)
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)
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}")
133class SyncMetricsCollector:
134 """同步性能指标收集器"""
136 def __init__(self):
137 self._metrics = SyncMetrics()
138 self._latencies: deque = deque(maxlen=1000)
139 self._lock = asyncio.Lock()
141 async def record_sync_start(self) -> None:
142 """记录同步开始"""
143 async with self._lock:
144 self._metrics.total_syncs += 1
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
155 # 更新平均延迟
156 self._latencies.append(duration_ms)
157 self._metrics.avg_latency_ms = sum(self._latencies) / len(self._latencies)
159 async def record_sync_error(self) -> None:
160 """记录同步错误"""
161 async with self._lock:
162 self._metrics.failed_syncs += 1
164 async def record_conflict(self) -> None:
165 """记录冲突"""
166 async with self._lock:
167 self._metrics.conflicts_resolved += 1
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 }
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
191 async def reset(self) -> None:
192 """重置指标"""
193 async with self._lock:
194 self._metrics = SyncMetrics()
195 self._latencies.clear()
198class SyncLogger:
199 """同步日志记录器"""
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()
206 async def log_event(self, event: SyncEvent) -> None:
207 """记录事件"""
208 async with self._lock:
209 self._events.append(event)
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"
218 logger.info(log_msg)
220 # 写入文件
221 if self._log_file:
222 await self._write_to_file(event)
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}")
239 async def get_recent_events(self, count: int = 50) -> List[SyncEvent]:
240 """获取最近事件"""
241 async with self._lock:
242 return list(self._events)[-count:]
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]
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]
254 async def clear(self) -> None:
255 """清空日志"""
256 async with self._lock:
257 self._events.clear()
260class SyncMonitor:
261 """同步监控器(整合模块)"""
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
269 async def start(self) -> None:
270 """启动监控"""
271 self._running = True
272 logger.info("同步监控器已启动")
274 async def stop(self) -> None:
275 """停止监控"""
276 self._running = False
277 logger.info("同步监控器已停止")
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()
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)
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)
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)
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()
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)
320 async def on_conflict(self, peer_id: str, conflict_type: str, resolution: str) -> None:
321 """冲突解决回调"""
322 await self.metrics_collector.record_conflict()
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)
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 }