zelos.observability
Phase 2 Observability — Structured Logging, Metrics, Tracing.
OpenTelemetry-compatible format. Prometheus export for metrics.
1""" 2Phase 2 Observability — Structured Logging, Metrics, Tracing. 3 4OpenTelemetry-compatible format. Prometheus export for metrics. 5""" 6 7import json 8import math 9import threading 10import time 11from collections.abc import Callable 12from dataclasses import dataclass, field 13from typing import Any, Optional 14 15# ═══════════════════ Structured Logging ═══════════════════ 16 17 18class StructuredLogger: 19 """JSON-formatted structured logger with level filtering.""" 20 21 LEVELS = {"debug": 10, "info": 20, "warn": 30, "error": 40} 22 23 def __init__(self, level: str = "info", format: str = "json"): 24 self.level = level 25 self.format = format 26 self._min_level = self.LEVELS.get(level, 20) 27 self._handlers: list[Callable] = [] 28 29 def add_handler(self, handler: Callable[[dict], None]) -> None: 30 self._handlers.append(handler) 31 32 def _log(self, level: str, message: str, **context) -> str | None: 33 if self.LEVELS.get(level, 0) < self._min_level: 34 return None 35 36 entry = { 37 "timestamp": time.time(), 38 "level": level, 39 "message": message, 40 "context": context, 41 } 42 43 if self.format == "json": 44 line = json.dumps(entry) 45 else: 46 line = f"[{level.upper()}] {message} {json.dumps(context) if context else ''}" 47 48 for h in self._handlers: 49 h(entry) 50 51 return line 52 53 def debug(self, message: str, **ctx) -> str | None: 54 return self._log("debug", message, **ctx) 55 56 def info(self, message: str, **ctx) -> str | None: 57 return self._log("info", message, **ctx) 58 59 def warn(self, message: str, **ctx) -> str | None: 60 return self._log("warn", message, **ctx) 61 62 def error(self, message: str, **ctx) -> str | None: 63 return self._log("error", message, **ctx) 64 65 66# ═══════════════════ Metrics ═══════════════════ 67 68 69class Counter: 70 """Monotonically increasing counter.""" 71 72 def __init__(self, name: str, help: str = "", labels: dict | None = None): 73 self.name = name 74 self.help = help 75 self._value: float = 0 76 self._labels = labels or {} 77 78 def inc(self, amount: float = 1) -> None: 79 self._value += amount 80 81 @property 82 def value(self) -> float: 83 return self._value 84 85 86class Gauge: 87 """Value that can go up and down.""" 88 89 def __init__(self, name: str, help: str = "", labels: dict | None = None): 90 self.name = name 91 self.help = help 92 self._value: float = 0 93 self._labels = labels or {} 94 95 def set(self, value: float) -> None: 96 self._value = value 97 98 def inc(self, amount: float = 1) -> None: 99 self._value += amount 100 101 def dec(self, amount: float = 1) -> None: 102 self._value -= amount 103 104 @property 105 def value(self) -> float: 106 return self._value 107 108 109class Histogram: 110 """Distribution of values with percentile calculation.""" 111 112 def __init__(self, name: str, help: str = "", buckets: list[float] | None = None): 113 self.name = name 114 self.help = help 115 self.buckets = buckets or [0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0] 116 self._values: list[float] = [] 117 self._lock = threading.Lock() 118 119 def observe(self, value: float) -> None: 120 with self._lock: 121 self._values.append(value) 122 123 def percentile(self, p: float) -> float: 124 with self._lock: 125 if not self._values: 126 return 0.0 127 sorted_vals = sorted(self._values) 128 idx = int(math.ceil(p / 100.0 * len(sorted_vals))) - 1 129 return sorted_vals[max(0, idx)] 130 131 @property 132 def count(self) -> int: 133 return len(self._values) 134 135 @property 136 def sum(self) -> float: 137 return sum(self._values) 138 139 140class MetricsCollector: 141 """Central metrics registry with Prometheus export support.""" 142 143 def __init__(self): 144 self._counters: dict[str, Counter] = {} 145 self._gauges: dict[str, Gauge] = {} 146 self._histograms: dict[str, Histogram] = {} 147 self._lock = threading.Lock() 148 149 def counter(self, name: str, help: str = "") -> Counter: 150 with self._lock: 151 if name not in self._counters: 152 self._counters[name] = Counter(name, help) 153 return self._counters[name] 154 155 def gauge(self, name: str, help: str = "") -> Gauge: 156 with self._lock: 157 if name not in self._gauges: 158 self._gauges[name] = Gauge(name, help) 159 return self._gauges[name] 160 161 def histogram(self, name: str, help: str = "") -> Histogram: 162 with self._lock: 163 if name not in self._histograms: 164 self._histograms[name] = Histogram(name, help) 165 return self._histograms[name] 166 167 def export_prometheus(self) -> str: 168 """Export all metrics in Prometheus text format.""" 169 lines = [] 170 for c in self._counters.values(): 171 lbl_str = ",".join(f'{k}="{v}"' for k, v in c._labels.items()) 172 name = f"{c.name}{{{lbl_str}}}" if lbl_str else c.name 173 lines.append(f"# HELP {c.name} {c.help}") 174 lines.append(f"# TYPE {c.name} counter") 175 lines.append(f"{name} {c.value}") 176 for g in self._gauges.values(): 177 lbl_str = ",".join(f'{k}="{v}"' for k, v in g._labels.items()) 178 name = f"{g.name}{{{lbl_str}}}" if lbl_str else g.name 179 lines.append(f"# HELP {g.name} {g.help}") 180 lines.append(f"# TYPE {g.name} gauge") 181 lines.append(f"{name} {g.value}") 182 for h in self._histograms.values(): 183 lines.append(f"# HELP {h.name} {h.help}") 184 lines.append(f"# TYPE {h.name} histogram") 185 lines.append(f"{h.name}_count {h.count}") 186 lines.append(f"{h.name}_sum {h.sum:.3f}") 187 return "\n".join(lines) + "\n" 188 189 def get_all(self) -> dict: 190 return { 191 "counters": {n: c.value for n, c in self._counters.items()}, 192 "gauges": {n: g.value for n, g in self._gauges.items()}, 193 "histograms": { 194 n: {"count": h.count, "p50": h.percentile(50), "p95": h.percentile(95), "p99": h.percentile(99)} 195 for n, h in self._histograms.items() 196 }, 197 } 198 199 200# ═══════════════════ Tracing ═══════════════════ 201 202 203@dataclass 204class SpanEvent: 205 name: str 206 timestamp: float = 0.0 207 attributes: dict[str, Any] = field(default_factory=dict) 208 209 210class Span: 211 def __init__(self, name: str, parent: Optional["Span"] = None): 212 self.name = name 213 self.span_id = str(hash(f"{name}{time.time()}"))[-16:] 214 self.parent_id = parent.span_id if parent else None 215 self.start_time = time.time() 216 self.end_time: float | None = None 217 self.events: list[SpanEvent] = [] 218 self.attributes: dict[str, Any] = {} 219 220 def add_event(self, name: str, **attributes) -> None: 221 self.events.append(SpanEvent(name, time.time(), attributes)) 222 223 def set_attribute(self, key: str, value: Any) -> None: 224 self.attributes[key] = value 225 226 def end(self) -> None: 227 self.end_time = time.time() 228 229 @property 230 def duration_ms(self) -> float: 231 if self.end_time: 232 return (self.end_time - self.start_time) * 1000 233 return 0 234 235 236class Tracer: 237 """Simple tracer with span hierarchy.""" 238 239 def __init__(self): 240 self._spans: list[Span] = [] 241 self._current: Span | None = None 242 self._stack: list[Span] = [] 243 244 def start_span(self, name: str) -> Span: 245 parent = self._stack[-1] if self._stack else None 246 span = Span(name, parent) 247 self._spans.append(span) 248 self._stack.append(span) 249 self._current = span 250 return span 251 252 def end_span(self) -> Span | None: 253 if self._stack: 254 span = self._stack.pop() 255 span.end() 256 self._current = self._stack[-1] if self._stack else None 257 return span 258 return None 259 260 def get_spans(self) -> list[Span]: 261 return list(self._spans)
class
StructuredLogger:
19class StructuredLogger: 20 """JSON-formatted structured logger with level filtering.""" 21 22 LEVELS = {"debug": 10, "info": 20, "warn": 30, "error": 40} 23 24 def __init__(self, level: str = "info", format: str = "json"): 25 self.level = level 26 self.format = format 27 self._min_level = self.LEVELS.get(level, 20) 28 self._handlers: list[Callable] = [] 29 30 def add_handler(self, handler: Callable[[dict], None]) -> None: 31 self._handlers.append(handler) 32 33 def _log(self, level: str, message: str, **context) -> str | None: 34 if self.LEVELS.get(level, 0) < self._min_level: 35 return None 36 37 entry = { 38 "timestamp": time.time(), 39 "level": level, 40 "message": message, 41 "context": context, 42 } 43 44 if self.format == "json": 45 line = json.dumps(entry) 46 else: 47 line = f"[{level.upper()}] {message} {json.dumps(context) if context else ''}" 48 49 for h in self._handlers: 50 h(entry) 51 52 return line 53 54 def debug(self, message: str, **ctx) -> str | None: 55 return self._log("debug", message, **ctx) 56 57 def info(self, message: str, **ctx) -> str | None: 58 return self._log("info", message, **ctx) 59 60 def warn(self, message: str, **ctx) -> str | None: 61 return self._log("warn", message, **ctx) 62 63 def error(self, message: str, **ctx) -> str | None: 64 return self._log("error", message, **ctx)
JSON-formatted structured logger with level filtering.
class
Counter:
70class Counter: 71 """Monotonically increasing counter.""" 72 73 def __init__(self, name: str, help: str = "", labels: dict | None = None): 74 self.name = name 75 self.help = help 76 self._value: float = 0 77 self._labels = labels or {} 78 79 def inc(self, amount: float = 1) -> None: 80 self._value += amount 81 82 @property 83 def value(self) -> float: 84 return self._value
Monotonically increasing counter.
class
Gauge:
87class Gauge: 88 """Value that can go up and down.""" 89 90 def __init__(self, name: str, help: str = "", labels: dict | None = None): 91 self.name = name 92 self.help = help 93 self._value: float = 0 94 self._labels = labels or {} 95 96 def set(self, value: float) -> None: 97 self._value = value 98 99 def inc(self, amount: float = 1) -> None: 100 self._value += amount 101 102 def dec(self, amount: float = 1) -> None: 103 self._value -= amount 104 105 @property 106 def value(self) -> float: 107 return self._value
Value that can go up and down.
class
Histogram:
110class Histogram: 111 """Distribution of values with percentile calculation.""" 112 113 def __init__(self, name: str, help: str = "", buckets: list[float] | None = None): 114 self.name = name 115 self.help = help 116 self.buckets = buckets or [0.01, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0] 117 self._values: list[float] = [] 118 self._lock = threading.Lock() 119 120 def observe(self, value: float) -> None: 121 with self._lock: 122 self._values.append(value) 123 124 def percentile(self, p: float) -> float: 125 with self._lock: 126 if not self._values: 127 return 0.0 128 sorted_vals = sorted(self._values) 129 idx = int(math.ceil(p / 100.0 * len(sorted_vals))) - 1 130 return sorted_vals[max(0, idx)] 131 132 @property 133 def count(self) -> int: 134 return len(self._values) 135 136 @property 137 def sum(self) -> float: 138 return sum(self._values)
Distribution of values with percentile calculation.
class
MetricsCollector:
141class MetricsCollector: 142 """Central metrics registry with Prometheus export support.""" 143 144 def __init__(self): 145 self._counters: dict[str, Counter] = {} 146 self._gauges: dict[str, Gauge] = {} 147 self._histograms: dict[str, Histogram] = {} 148 self._lock = threading.Lock() 149 150 def counter(self, name: str, help: str = "") -> Counter: 151 with self._lock: 152 if name not in self._counters: 153 self._counters[name] = Counter(name, help) 154 return self._counters[name] 155 156 def gauge(self, name: str, help: str = "") -> Gauge: 157 with self._lock: 158 if name not in self._gauges: 159 self._gauges[name] = Gauge(name, help) 160 return self._gauges[name] 161 162 def histogram(self, name: str, help: str = "") -> Histogram: 163 with self._lock: 164 if name not in self._histograms: 165 self._histograms[name] = Histogram(name, help) 166 return self._histograms[name] 167 168 def export_prometheus(self) -> str: 169 """Export all metrics in Prometheus text format.""" 170 lines = [] 171 for c in self._counters.values(): 172 lbl_str = ",".join(f'{k}="{v}"' for k, v in c._labels.items()) 173 name = f"{c.name}{{{lbl_str}}}" if lbl_str else c.name 174 lines.append(f"# HELP {c.name} {c.help}") 175 lines.append(f"# TYPE {c.name} counter") 176 lines.append(f"{name} {c.value}") 177 for g in self._gauges.values(): 178 lbl_str = ",".join(f'{k}="{v}"' for k, v in g._labels.items()) 179 name = f"{g.name}{{{lbl_str}}}" if lbl_str else g.name 180 lines.append(f"# HELP {g.name} {g.help}") 181 lines.append(f"# TYPE {g.name} gauge") 182 lines.append(f"{name} {g.value}") 183 for h in self._histograms.values(): 184 lines.append(f"# HELP {h.name} {h.help}") 185 lines.append(f"# TYPE {h.name} histogram") 186 lines.append(f"{h.name}_count {h.count}") 187 lines.append(f"{h.name}_sum {h.sum:.3f}") 188 return "\n".join(lines) + "\n" 189 190 def get_all(self) -> dict: 191 return { 192 "counters": {n: c.value for n, c in self._counters.items()}, 193 "gauges": {n: g.value for n, g in self._gauges.items()}, 194 "histograms": { 195 n: {"count": h.count, "p50": h.percentile(50), "p95": h.percentile(95), "p99": h.percentile(99)} 196 for n, h in self._histograms.items() 197 }, 198 }
Central metrics registry with Prometheus export support.
def
export_prometheus(self) -> str:
168 def export_prometheus(self) -> str: 169 """Export all metrics in Prometheus text format.""" 170 lines = [] 171 for c in self._counters.values(): 172 lbl_str = ",".join(f'{k}="{v}"' for k, v in c._labels.items()) 173 name = f"{c.name}{{{lbl_str}}}" if lbl_str else c.name 174 lines.append(f"# HELP {c.name} {c.help}") 175 lines.append(f"# TYPE {c.name} counter") 176 lines.append(f"{name} {c.value}") 177 for g in self._gauges.values(): 178 lbl_str = ",".join(f'{k}="{v}"' for k, v in g._labels.items()) 179 name = f"{g.name}{{{lbl_str}}}" if lbl_str else g.name 180 lines.append(f"# HELP {g.name} {g.help}") 181 lines.append(f"# TYPE {g.name} gauge") 182 lines.append(f"{name} {g.value}") 183 for h in self._histograms.values(): 184 lines.append(f"# HELP {h.name} {h.help}") 185 lines.append(f"# TYPE {h.name} histogram") 186 lines.append(f"{h.name}_count {h.count}") 187 lines.append(f"{h.name}_sum {h.sum:.3f}") 188 return "\n".join(lines) + "\n"
Export all metrics in Prometheus text format.
def
get_all(self) -> dict:
190 def get_all(self) -> dict: 191 return { 192 "counters": {n: c.value for n, c in self._counters.items()}, 193 "gauges": {n: g.value for n, g in self._gauges.items()}, 194 "histograms": { 195 n: {"count": h.count, "p50": h.percentile(50), "p95": h.percentile(95), "p99": h.percentile(99)} 196 for n, h in self._histograms.items() 197 }, 198 }
@dataclass
class
SpanEvent:
class
Span:
211class Span: 212 def __init__(self, name: str, parent: Optional["Span"] = None): 213 self.name = name 214 self.span_id = str(hash(f"{name}{time.time()}"))[-16:] 215 self.parent_id = parent.span_id if parent else None 216 self.start_time = time.time() 217 self.end_time: float | None = None 218 self.events: list[SpanEvent] = [] 219 self.attributes: dict[str, Any] = {} 220 221 def add_event(self, name: str, **attributes) -> None: 222 self.events.append(SpanEvent(name, time.time(), attributes)) 223 224 def set_attribute(self, key: str, value: Any) -> None: 225 self.attributes[key] = value 226 227 def end(self) -> None: 228 self.end_time = time.time() 229 230 @property 231 def duration_ms(self) -> float: 232 if self.end_time: 233 return (self.end_time - self.start_time) * 1000 234 return 0
Span(name: str, parent: Optional[Span] = None)
212 def __init__(self, name: str, parent: Optional["Span"] = None): 213 self.name = name 214 self.span_id = str(hash(f"{name}{time.time()}"))[-16:] 215 self.parent_id = parent.span_id if parent else None 216 self.start_time = time.time() 217 self.end_time: float | None = None 218 self.events: list[SpanEvent] = [] 219 self.attributes: dict[str, Any] = {}
events: list[SpanEvent]
class
Tracer:
237class Tracer: 238 """Simple tracer with span hierarchy.""" 239 240 def __init__(self): 241 self._spans: list[Span] = [] 242 self._current: Span | None = None 243 self._stack: list[Span] = [] 244 245 def start_span(self, name: str) -> Span: 246 parent = self._stack[-1] if self._stack else None 247 span = Span(name, parent) 248 self._spans.append(span) 249 self._stack.append(span) 250 self._current = span 251 return span 252 253 def end_span(self) -> Span | None: 254 if self._stack: 255 span = self._stack.pop() 256 span.end() 257 self._current = self._stack[-1] if self._stack else None 258 return span 259 return None 260 261 def get_spans(self) -> list[Span]: 262 return list(self._spans)
Simple tracer with span hierarchy.