1#!/usr/bin/env python3
2"""Check live metric cache continuity when its upstream stops answering."""
3import argparse
4from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
5import json
6from pathlib import Path
7import shlex
8import subprocess
9import tempfile
10import threading
11import time
12import urllib.error
13import urllib.request
14
15parser = argparse.ArgumentParser(description=__doc__)
16parser.add_argument("binary")
17args = parser.parse_args()
18failed = False
19reads = 0
20
21class Metrics(BaseHTTPRequestHandler):
22 def log_message(self, *args):
23 pass
24
25 def do_GET(self):
26 global reads
27 reads += 1
28 self.send_response(503 if failed else 200)
29 self.send_header("Content-Type", "application/json")
30 self.end_headers()
31 self.wfile.write(json.dumps({"status": "success", "data": {"result": [
32 {"metric": {}, "values": [[time.time(), "12.5"]]}
33 ]}}).encode())
34
35metrics = ThreadingHTTPServer(("127.0.0.1", 0), Metrics)
36
37class Nomad(BaseHTTPRequestHandler):
38 def log_message(self, *args):
39 pass
40
41 def do_GET(self):
42 self.send_response(200)
43 self.send_header("Content-Type", "application/json")
44 self.end_headers()
45 self.wfile.write(json.dumps([
46 {"Address": "127.0.0.1", "Port": metrics.server_port}
47 ] if self.path == "/v1/service/victoria-metrics" else []).encode())
48
49nomad = ThreadingHTTPServer(("127.0.0.1", 0), Nomad)
50for server in [metrics, nomad]:
51 threading.Thread(target=server.serve_forever, daemon=True).start()
52
53unit = "studio-dashboard-metrics-test"
54with tempfile.TemporaryDirectory(prefix="studio-metrics-test-") as temporary:
55 environment = subprocess.check_output(["systemctl", "show", "-P", "Environment", "studio-dashboard"], text=True)
56 variables = dict(item.split("=", 1) for item in shlex.split(environment))
57 for key in ["STUDIO_INDEX_POOL", "STUDIO_YT_STATE", "STUDIO_FILES_WRITABLE"]:
58 variables.pop(key, None)
59 variables.update(PORT="7074", STUDIO_METRICS="follow", STUDIO_DATA_DIR=temporary,
60 NOMAD_ADDR=f"http://127.0.0.1:{nomad.server_port}")
61 proof = {"Studio-Proxy-Token": Path(variables["STUDIO_PROXY_TOKEN_FILE"]).read_text().strip()} if "STUDIO_PROXY_TOKEN_FILE" in variables else {}
62 def request(range=3600):
63 started = time.monotonic()
64 req = urllib.request.Request(f"http://127.0.0.1:7074/api/metrics/host.cpu?range={range}",
65 headers={"User-Name": "fixture", "User-Groups": "infra-admin", **proof})
66 with urllib.request.urlopen(req, timeout=10) as response:
67 return json.load(response), time.monotonic() - started
68
69 try:
70 subprocess.run(["systemd-run", "--unit=" + unit, "--collect", "--property=CPUWeight=1000",
71 *["--setenv=" + k + "=" + v for k, v in variables.items()],
72 str(Path(args.binary).resolve())], check=True, capture_output=True)
73 for _ in range(100):
74 try:
75 baseline, _ = request()
76 break
77 except OSError:
78 time.sleep(.1)
79 else:
80 raise AssertionError("fixture did not start")
81 assert baseline[0]["v"] == [12.5], baseline
82 failed = True
83 time.sleep(2.2)
84 cached, elapsed = request()
85 assert cached == baseline, (cached, baseline)
86 assert elapsed < 2, elapsed
87 time.sleep(2.2)
88 cached, elapsed = request()
89 assert cached == baseline, (cached, baseline)
90 try:
91 request(7200)
92 except urllib.error.HTTPError as error:
93 assert error.code == 502, error.code
94 else:
95 raise AssertionError("a different range reused the cached series")
96 failed = False
97 for _ in range(30):
98 cached, _ = request()
99 if cached != baseline:
100 break
101 time.sleep(.2)
102 else:
103 raise AssertionError("the live series did not refresh after recovery")
104 print(json.dumps({"stale_live_series": "passed", "range_isolation": "passed",
105 "recovery_refresh": "passed", "last_cached_request_ms": round(elapsed * 1000, 1),
106 "upstream_reads": reads}))
107 finally:
108 subprocess.run(["systemctl", "stop", unit], capture_output=True)
109 nomad.shutdown()
110 metrics.shutdown()