-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathaggregator.py
More file actions
143 lines (123 loc) · 5.69 KB
/
Copy pathaggregator.py
File metadata and controls
143 lines (123 loc) · 5.69 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
import json
import threading
import time
import math
def _percentile(sorted_vals: list[float], p: float) -> float:
if not sorted_vals:
return 0.0
idx = max(0, min(math.ceil(p * len(sorted_vals)) - 1, len(sorted_vals) - 1))
return sorted_vals[idx]
class Aggregator:
def __init__(self, transport, flush_interval_ms: int = 60_000):
self._transport = transport
self._flush_interval = flush_interval_ms / 1000.0
self._buffer: dict = {}
self._lock = threading.Lock()
self._timer: threading.Timer | None = None
def start(self):
self._schedule()
def stop(self):
if self._timer:
self._timer.cancel()
self._timer = None
self._flush()
def record(self, event: dict):
is_ghost = event.get("is_ghost", False)
key = f"{event['method']}|{event['route']}|{event['env']}|{event.get('release') or ''}|{'1' if is_ghost else '0'}"
with self._lock:
if key not in self._buffer:
self._buffer[key] = {
"method": event["method"],
"route": event["route"],
"env": event["env"],
"release": event.get("release"),
"is_ghost": is_ghost,
"durations": [],
"ttfb_durations": [],
"response_sizes": [],
"request_sizes": [],
"inflight_samples": [],
"status_2xx": 0,
"status_3xx": 0,
"status_4xx": 0,
"status_5xx": 0,
"status_map": {},
}
bucket = self._buffer[key]
bucket["durations"].append(event["duration_ms"])
if event.get("ttfb_ms") is not None:
bucket["ttfb_durations"].append(event["ttfb_ms"])
if event.get("response_size") is not None:
bucket["response_sizes"].append(event["response_size"])
if event.get("request_size") is not None:
bucket["request_sizes"].append(event["request_size"])
if event.get("inflight") is not None:
bucket["inflight_samples"].append(event["inflight"])
s = event["status"]
if 200 <= s < 300: bucket["status_2xx"] += 1
elif 300 <= s < 400: bucket["status_3xx"] += 1
elif 400 <= s < 500: bucket["status_4xx"] += 1
elif s >= 500: bucket["status_5xx"] += 1
bucket["status_map"][s] = bucket["status_map"].get(s, 0) + 1
def _schedule(self):
self._timer = threading.Timer(self._flush_interval, self._tick)
self._timer.daemon = True
self._timer.start()
def _tick(self):
self._flush()
self._schedule()
def _flush(self):
with self._lock:
if not self._buffer:
return
snapshot = self._buffer
self._buffer = {}
bucket_ts = int(time.time() // 60) * 60
rows = []
for bucket in snapshot.values():
sorted_d = sorted(bucket["durations"])
sorted_ttfb = sorted(bucket["ttfb_durations"])
n = len(sorted_d)
sizes = bucket["response_sizes"]
req_sizes = bucket["request_sizes"]
inflight = bucket["inflight_samples"]
bytes_avg = sum(sizes) / len(sizes) if sizes else None
request_size_avg = sum(req_sizes) / len(req_sizes) if req_sizes else None
lat_avg = sum(bucket["durations"]) / n if n > 0 else None
inflight_avg = sum(inflight) / len(inflight) if inflight else None
inflight_max = max(inflight) if inflight else None
# Granular distribution — sorted by count desc
status_dist = (
json.dumps(
dict(sorted(bucket["status_map"].items(), key=lambda x: -x[1]))
)
if bucket["status_map"] else None
)
rows.append({
"bucket_ts": bucket_ts,
"route": bucket["route"],
"method": bucket["method"],
"env": bucket["env"],
"release_tag": bucket["release"],
"is_ghost": 1 if bucket["is_ghost"] else 0,
"status_2xx": bucket["status_2xx"],
"status_3xx": bucket["status_3xx"],
"status_4xx": bucket["status_4xx"],
"status_5xx": bucket["status_5xx"],
"status_dist": status_dist,
"total_calls": n,
"lat_p50": _percentile(sorted_d, 0.50),
"lat_p90": _percentile(sorted_d, 0.90),
"lat_p99": _percentile(sorted_d, 0.99),
"lat_avg": lat_avg,
"lat_min": sorted_d[0] if sorted_d else 0,
"lat_max": sorted_d[-1] if sorted_d else 0,
"lat_ttfb_p50": _percentile(sorted_ttfb, 0.50) if sorted_ttfb else None,
"lat_ttfb_p90": _percentile(sorted_ttfb, 0.90) if sorted_ttfb else None,
"lat_ttfb_p99": _percentile(sorted_ttfb, 0.99) if sorted_ttfb else None,
"bytes_avg": bytes_avg,
"request_size_avg": request_size_avg,
"inflight_avg": inflight_avg,
"inflight_max": inflight_max,
})
self._transport.write(rows)