jetstream v2 in zig stream.waow.tech
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139#!/usr/bin/env python3"""Empirical contract for Stream's real process metrics endpoint."""
from __future__ import annotations
import mathimport osimport reimport timeimport urllib.request
BASE = os.environ.get("STREAM_BASE_URL", "http://127.0.0.1:6008")TYPE_RE = re.compile(r"^# TYPE (\S+) (\S+)$")SAMPLE_RE = re.compile(r"^(\S+?)(?:\{[^}]*\})?\s+([^\s]+)$")STORE_COUNT_RE = re.compile( r'^jetstream_store_op_duration_seconds_count\{op="([^"]+)",status="([^"]+)"\}\s+([^\s]+)$', re.MULTILINE,)STORE_BUCKET_RE = re.compile( r'^jetstream_store_op_duration_seconds_bucket\{op="([^"]+)",status="([^"]+)",le="([^"]+)"\}\s+([^\s]+)$', re.MULTILINE,)
EXPECTED_TYPES = { "process_start_time_seconds": "gauge", "process_cpu_seconds_total": "counter", "process_resident_memory_bytes": "gauge", "process_virtual_memory_bytes": "gauge", "stream_process_threads": "gauge", "process_open_fds": "gauge", "process_max_fds": "gauge", "process_minor_page_faults_total": "counter", "process_major_page_faults_total": "counter", "process_voluntary_context_switches_total": "counter", "process_involuntary_context_switches_total": "counter",}COUNTERS = {name for name, kind in EXPECTED_TYPES.items() if kind == "counter"}
def scrape() -> tuple[str, dict[str, str], dict[str, float]]: with urllib.request.urlopen(BASE + "/metrics", timeout=5) as response: assert response.status == 200 raw = response.read().decode() types: dict[str, str] = {} samples: dict[str, float] = {} for line in raw.splitlines(): if match := TYPE_RE.match(line): types[match.group(1)] = match.group(2) elif not line.startswith("#") and (match := SAMPLE_RE.match(line)): samples[match.group(1)] = float(match.group(2)) return raw, types, samples
def wait_past_startup_gate() -> None: """listener-first startup: /healthz and /metrics answer from bind time, before the storage reads this contract counts (the fresh-dir phase read is the guaranteed get-notfound). /status leaves "starting" only after storage is published.""" deadline = time.monotonic() + 15 while time.monotonic() < deadline: try: with urllib.request.urlopen(BASE + "/status", timeout=1) as response: if b"state starting" not in response.read(): return except OSError: pass time.sleep(0.05) raise AssertionError("startup gate never opened")
def main() -> None: wait_past_startup_gate() first_raw, first_types, first = scrape() assert not re.search(r"(?m)^go_", first_raw), "Stream emitted synthetic Go metrics" for name, kind in EXPECTED_TYPES.items(): assert first_types.get(name) == kind, (name, first_types.get(name), kind) assert name in first and math.isfinite(first[name]) and first[name] >= 0, (name, first.get(name))
assert first_types.get("jetstream_store_op_duration_seconds") == "histogram" store_counts = { (op, status): float(value) for op, status, value in STORE_COUNT_RE.findall(first_raw) } expected_store_series = { ("get", "ok"), ("get", "notfound"), ("get", "error"), ("set", "ok"), ("set", "error"), ("delete", "ok"), ("delete", "error"), ("batch_commit", "ok"), ("batch_commit", "error"), } assert set(store_counts) == expected_store_series, store_counts assert store_counts[("get", "notfound")] > 0 assert store_counts[("set", "ok")] > 0 assert store_counts[("batch_commit", "ok")] > 0 buckets: dict[tuple[str, str], list[tuple[str, float]]] = {} for op, status, upper, value in STORE_BUCKET_RE.findall(first_raw): buckets.setdefault((op, status), []).append((upper, float(value))) assert set(buckets) == expected_store_series for labels, observations in buckets.items(): assert len(observations) == 16, (labels, observations) assert observations[-1] == ("+Inf", store_counts[labels]), (labels, observations[-1]) values = [value for _, value in observations] assert values == sorted(values), (labels, values)
assert 0 < first["process_start_time_seconds"] <= time.time() assert 0 < first["process_resident_memory_bytes"] <= first["process_virtual_memory_bytes"] assert first["stream_process_threads"] >= 1 assert 0 < first["process_open_fds"] <= first["process_max_fds"]
# Exercise the handler between snapshots, then prove cumulative process # counters never moved backwards and identity stayed fixed. for _ in range(8): urllib.request.urlopen(BASE + "/healthz", timeout=5).read() _, second_types, second = scrape() assert second_types == first_types assert second["process_start_time_seconds"] == first["process_start_time_seconds"] for name in COUNTERS: assert second[name] >= first[name], (name, first[name], second[name])
print( "process metrics PASS: " f"rss={int(second['process_resident_memory_bytes'])} " f"vm={int(second['process_virtual_memory_bytes'])} " f"threads={int(second['stream_process_threads'])} " f"fds={int(second['process_open_fds'])}/{int(second['process_max_fds'])} " f"store=get-notfound:{int(store_counts[('get', 'notfound')])}," f"set:{int(store_counts[('set', 'ok')])}," f"batch:{int(store_counts[('batch_commit', 'ok')])}" )
if __name__ == "__main__": main()