| 1 | #!/usr/bin/env python3 |
| 2 | """Check live metric cache continuity when its upstream stops answering.""" |
| 3 | import argparse |
| 4 | from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer |
| 5 | import json |
| 6 | from pathlib import Path |
| 7 | import shlex |
| 8 | import subprocess |
| 9 | import tempfile |
| 10 | import threading |
| 11 | import time |
| 12 | import urllib.error |
| 13 | import urllib.request |
| 14 | |
| 15 | parser = argparse.ArgumentParser(description=__doc__) |
| 16 | parser.add_argument("binary") |
| 17 | args = parser.parse_args() |
| 18 | failed = False |
| 19 | reads = 0 |
| 20 | |
| 21 | class 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 | |
| 35 | metrics = ThreadingHTTPServer(("127.0.0.1", 0), Metrics) |
| 36 | |
| 37 | class 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 | |
| 49 | nomad = ThreadingHTTPServer(("127.0.0.1", 0), Nomad) |
| 50 | for server in [metrics, nomad]: |
| 51 | threading.Thread(target=server.serve_forever, daemon=True).start() |
| 52 | |
| 53 | unit = "studio-dashboard-metrics-test" |
| 54 | with 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() |