Coverage for agentos/observability/metrics.py: 48%

169 statements  

« prev     ^ index     » next       coverage.py v7.14.3, created at 2026-07-09 10:19 +0800

1""" # noqa: E501 

2AgentOS v0.70 — 性能指标与可观测性增强。 

3基因来源: Prometheus metrics + OpenTelemetry 

4 

5提供: 

6- 延迟分位数 (p50/p95/p99) 

7- 吞吐量统计 (RPS) 

8- 错误率追踪 

9- 缓存命中率 

10- TTL-based环形缓冲区 

11""" 

12 

13from __future__ import annotations 

14 

15import threading 

16import time 

17from collections import deque 

18from dataclasses import dataclass, field 

19 

20 

21@dataclass 

22class MetricSnapshot: 

23 """指标快照 — 用于导出/序列化。""" 

24 

25 timestamp: float = field(default_factory=time.time) 

26 histograms: dict[str, dict] = field(default_factory=dict) 

27 counters: dict[str, int] = field(default_factory=dict) 

28 gauges: dict[str, float] = field(default_factory=dict) 

29 derived_metrics: dict[str, float] = field(default_factory=dict) 

30 

31 def to_json(self) -> str: 

32 import json 

33 

34 return json.dumps( 

35 { 

36 "ts": self.timestamp, 

37 "h": self.histograms, 

38 "c": self.counters, 

39 "g": self.gauges, 

40 "d": self.derived_metrics, 

41 } 

42 ) 

43 

44 @classmethod 

45 def from_collector(cls, collector: MetricsCollector) -> MetricSnapshot: 

46 s = collector.snapshot() 

47 return cls( 

48 histograms={ 

49 "step": s["latency_step_ms"], 

50 "model": s["latency_model_ms"], 

51 "tool": s["latency_tool_ms"], 

52 }, 

53 counters={ 

54 "steps": collector.steps_total.value, 

55 "model_calls": collector.model_calls_total.value, 

56 "tool_calls": collector.tool_calls_total.value, 

57 "errors": collector.errors_total.value, 

58 "cache_hits": collector.cache_hits.value, 

59 "cache_misses": collector.cache_misses.value, 

60 }, 

61 gauges={ 

62 "active_agents": collector.active_agents.value, 

63 "queue_depth": collector.queue_depth.value, 

64 }, 

65 derived_metrics={ 

66 "rps": s["throughput"]["rps"], 

67 "error_rate": s["error_rate"], 

68 "cache_hit_rate": s["cache_hit_rate"], 

69 }, 

70 ) 

71 

72 

73@dataclass 

74class MetricPoint: 

75 """指标数据点。""" 

76 

77 timestamp: float 

78 value: float 

79 labels: dict[str, str] = field(default_factory=dict) 

80 

81 

82@dataclass 

83class Histogram: 

84 """滑动窗口直方图 — 计算分位数。""" 

85 

86 name: str 

87 window_seconds: float = 300.0 

88 max_size: int = 10000 

89 _points: deque = field(default_factory=deque) 

90 _lock: threading.Lock = field(default_factory=threading.Lock) 

91 

92 def observe(self, value: float, **labels): 

93 with self._lock: 

94 self._points.append(MetricPoint(timestamp=time.time(), value=value, labels=labels)) 

95 self._prune() 

96 if len(self._points) > self.max_size: 

97 self._points.popleft() 

98 

99 def _prune(self): 

100 cutoff = time.time() - self.window_seconds 

101 while self._points and self._points[0].timestamp < cutoff: 

102 self._points.popleft() 

103 

104 @property 

105 def count(self) -> int: 

106 with self._lock: 

107 self._prune() 

108 return len(self._points) 

109 

110 def quantile(self, q: float) -> float: 

111 """计算分位数 0.5=p50, 0.95=p95, 0.99=p99。""" 

112 with self._lock: 

113 self._prune() 

114 if not self._points: 

115 return 0.0 

116 values = sorted(p.value for p in self._points) 

117 idx = int(len(values) * q) 

118 if idx >= len(values): 

119 idx = len(values) - 1 

120 return values[idx] 

121 

122 @property 

123 def p50(self) -> float: 

124 return self.quantile(0.5) 

125 

126 @property 

127 def p95(self) -> float: 

128 return self.quantile(0.95) 

129 

130 @property 

131 def p99(self) -> float: 

132 return self.quantile(0.99) 

133 

134 @property 

135 def avg(self) -> float: 

136 with self._lock: 

137 self._prune() 

138 if not self._points: 

139 return 0.0 

140 return sum(p.value for p in self._points) / len(self._points) 

141 

142 @property 

143 def min_val(self) -> float: 

144 with self._lock: 

145 self._prune() 

146 if not self._points: 

147 return 0.0 

148 return min(p.value for p in self._points) 

149 

150 @property 

151 def max_val(self) -> float: 

152 with self._lock: 

153 self._prune() 

154 if not self._points: 

155 return 0.0 

156 return max(p.value for p in self._points) 

157 

158 def stats(self) -> dict: 

159 return { 

160 "name": self.name, 

161 "count": self.count, 

162 "avg": self.avg, 

163 "p50": self.p50, 

164 "p95": self.p95, 

165 "p99": self.p99, 

166 "min": self.min_val, 

167 "max": self.max_val, 

168 "window_seconds": self.window_seconds, 

169 } 

170 

171 

172@dataclass 

173class Counter: 

174 """单调递增计数器。""" 

175 

176 name: str 

177 _value: int = 0 

178 _labels: dict[str, str] = field(default_factory=dict) 

179 

180 def inc(self, amount: int = 1): 

181 self._value += amount 

182 

183 @property 

184 def value(self) -> int: 

185 return self._value 

186 

187 

188@dataclass 

189class Gauge: 

190 """可增可减的仪表值。""" 

191 

192 name: str 

193 _value: float = 0.0 

194 _labels: dict[str, str] = field(default_factory=dict) 

195 

196 def set(self, value: float): 

197 self._value = value 

198 

199 def inc(self, amount: float = 1.0): 

200 self._value += amount 

201 

202 def dec(self, amount: float = 1.0): 

203 self._value -= amount 

204 

205 @property 

206 def value(self) -> float: 

207 return self._value 

208 

209 

210class MetricsCollector: 

211 """ 

212 统一指标收集器。 

213 内置: latency, throughput, error_rate, cache_hit_rate。 

214 """ 

215 

216 def __init__(self, window_seconds: float = 300.0): 

217 self.window_seconds = window_seconds 

218 

219 # Histograms 

220 self.latency_step = Histogram("step_latency", window_seconds) 

221 self.latency_model = Histogram("model_latency", window_seconds) 

222 self.latency_tool = Histogram("tool_latency", window_seconds) 

223 

224 # Counters 

225 self.steps_total = Counter("steps_total") 

226 self.model_calls_total = Counter("model_calls_total") 

227 self.tool_calls_total = Counter("tool_calls_total") 

228 self.errors_total = Counter("errors_total") 

229 self.cache_hits = Counter("cache_hits") 

230 self.cache_misses = Counter("cache_misses") 

231 

232 # Gauges 

233 self.active_agents = Gauge("active_agents") 

234 self.queue_depth = Gauge("queue_depth") 

235 self.memory_used_mb = Gauge("memory_used_mb") 

236 

237 self._start_time = time.time() 

238 

239 # ── Recording ──────────────────────────────── 

240 

241 def record_step_latency(self, duration_ms: float): 

242 self.latency_step.observe(duration_ms) 

243 self.steps_total.inc() 

244 

245 def record_model_latency(self, duration_ms: float, model: str = ""): 

246 self.latency_model.observe(duration_ms, model=model) 

247 self.model_calls_total.inc() 

248 

249 def record_tool_latency(self, duration_ms: float, tool: str = ""): 

250 self.latency_tool.observe(duration_ms, tool=tool) 

251 self.tool_calls_total.inc() 

252 

253 def record_error(self): 

254 self.errors_total.inc() 

255 

256 def record_cache_hit(self): 

257 self.cache_hits.inc() 

258 

259 def record_cache_miss(self): 

260 self.cache_misses.inc() 

261 

262 # ── Derived Metrics ────────────────────────── 

263 

264 @property 

265 def uptime_seconds(self) -> float: 

266 return time.time() - self._start_time 

267 

268 @property 

269 def rps(self) -> float: 

270 """请求速率 (steps/sec over window)。""" 

271 if self.uptime_seconds < 1: 

272 return self.steps_total.value 

273 return self.steps_total.value / self.uptime_seconds 

274 

275 @property 

276 def error_rate(self) -> float: 

277 total = self.steps_total.value + self.errors_total.value 

278 if total == 0: 

279 return 0.0 

280 return self.errors_total.value / total 

281 

282 @property 

283 def cache_hit_rate(self) -> float: 

284 total = self.cache_hits.value + self.cache_misses.value 

285 if total == 0: 

286 return 0.0 

287 return self.cache_hits.value / total 

288 

289 # ── Snapshot ───────────────────────────────── 

290 

291 def snapshot(self) -> dict: 

292 return { 

293 "uptime_seconds": self.uptime_seconds, 

294 "throughput": { 

295 "rps": round(self.rps, 2), 

296 "steps_total": self.steps_total.value, 

297 "model_calls": self.model_calls_total.value, 

298 "tool_calls": self.tool_calls_total.value, 

299 }, 

300 "latency_step_ms": self.latency_step.stats(), 

301 "latency_model_ms": self.latency_model.stats(), 

302 "latency_tool_ms": self.latency_tool.stats(), 

303 "error_rate": round(self.error_rate, 4), 

304 "errors_total": self.errors_total.value, 

305 "cache_hit_rate": round(self.cache_hit_rate, 2), 

306 "cache_hits": self.cache_hits.value, 

307 "cache_misses": self.cache_misses.value, 

308 "active_agents": self.active_agents.value, 

309 "queue_depth": self.queue_depth.value, 

310 } 

311 

312 def summary(self) -> str: 

313 s = self.snapshot() 

314 lines = [ 

315 f"运行时间: {s['uptime_seconds']:.0f}s", 

316 f"吞吐: {s['throughput']['rps']} rps ({s['throughput']['steps_total']} steps)", 

317 f"延迟: p50={s['latency_step_ms']['p50']:.0f}ms p95={s['latency_step_ms']['p95']:.0f}ms p99={s['latency_step_ms']['p99']:.0f}ms", # noqa: E501 

318 f"错误率: {s['error_rate']:.2%} ({s['errors_total']} errors)", 

319 f"缓存命中率: {s['cache_hit_rate']:.1%} ({s['cache_hits']}/{s['cache_hits'] + s['cache_misses']})", 

320 ] 

321 return "\n".join(lines)