Coverage for agentos/observability/metrics.py: 48%
169 statements
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-10 01:20 +0800
« prev ^ index » next coverage.py v7.14.3, created at 2026-07-10 01:20 +0800
1""" # noqa: E501
2AgentOS v0.70 — 性能指标与可观测性增强。
3基因来源: Prometheus metrics + OpenTelemetry
5提供:
6- 延迟分位数 (p50/p95/p99)
7- 吞吐量统计 (RPS)
8- 错误率追踪
9- 缓存命中率
10- TTL-based环形缓冲区
11"""
13from __future__ import annotations
15import threading
16import time
17from collections import deque
18from dataclasses import dataclass, field
21@dataclass
22class MetricSnapshot:
23 """指标快照 — 用于导出/序列化。"""
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)
31 def to_json(self) -> str:
32 import json
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 )
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 )
73@dataclass
74class MetricPoint:
75 """指标数据点。"""
77 timestamp: float
78 value: float
79 labels: dict[str, str] = field(default_factory=dict)
82@dataclass
83class Histogram:
84 """滑动窗口直方图 — 计算分位数。"""
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)
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()
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()
104 @property
105 def count(self) -> int:
106 with self._lock:
107 self._prune()
108 return len(self._points)
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]
122 @property
123 def p50(self) -> float:
124 return self.quantile(0.5)
126 @property
127 def p95(self) -> float:
128 return self.quantile(0.95)
130 @property
131 def p99(self) -> float:
132 return self.quantile(0.99)
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)
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)
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)
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 }
172@dataclass
173class Counter:
174 """单调递增计数器。"""
176 name: str
177 _value: int = 0
178 _labels: dict[str, str] = field(default_factory=dict)
180 def inc(self, amount: int = 1):
181 self._value += amount
183 @property
184 def value(self) -> int:
185 return self._value
188@dataclass
189class Gauge:
190 """可增可减的仪表值。"""
192 name: str
193 _value: float = 0.0
194 _labels: dict[str, str] = field(default_factory=dict)
196 def set(self, value: float):
197 self._value = value
199 def inc(self, amount: float = 1.0):
200 self._value += amount
202 def dec(self, amount: float = 1.0):
203 self._value -= amount
205 @property
206 def value(self) -> float:
207 return self._value
210class MetricsCollector:
211 """
212 统一指标收集器。
213 内置: latency, throughput, error_rate, cache_hit_rate。
214 """
216 def __init__(self, window_seconds: float = 300.0):
217 self.window_seconds = window_seconds
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)
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")
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")
237 self._start_time = time.time()
239 # ── Recording ────────────────────────────────
241 def record_step_latency(self, duration_ms: float):
242 self.latency_step.observe(duration_ms)
243 self.steps_total.inc()
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()
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()
253 def record_error(self):
254 self.errors_total.inc()
256 def record_cache_hit(self):
257 self.cache_hits.inc()
259 def record_cache_miss(self):
260 self.cache_misses.inc()
262 # ── Derived Metrics ──────────────────────────
264 @property
265 def uptime_seconds(self) -> float:
266 return time.time() - self._start_time
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
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
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
289 # ── Snapshot ─────────────────────────────────
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 }
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)