"""Observe configured compaction controls in a real Stream process.""" import json import os import pathlib import time import urllib.request port = int(os.environ["STREAM_COMPACTION_CONFIG_PORT"]) expected_chunks = int(os.environ["STREAM_COMPACTION_CONFIG_CHUNKS"]) expected_workers = int(os.environ["STREAM_COMPACTION_CONFIG_WORKERS"]) log_path = pathlib.Path(os.environ["STREAM_COMPACTION_CONFIG_LOG"]) base = f"http://127.0.0.1:{port}" deadline = time.monotonic() + 20 while time.monotonic() < deadline: log = log_path.read_text(errors="replace") if log_path.exists() else "" marker = f"compaction pass: watermark 4, {expected_chunks} chunks," workers = f", {expected_workers} rewrite workers" if marker in log and workers in log: listing = json.load( urllib.request.urlopen( base + "/xrpc/network.bsky.jetstream.listSegments", timeout=2 ) ) assert sum(segment["eventCount"] for segment in listing["segments"]) == 2 metrics = urllib.request.urlopen(base + "/metrics", timeout=2).read().decode() assert "jetstream_compaction_watermark_seq 4" in metrics print( "compaction config: PASS " f"(chunks={expected_chunks}, rewrite_workers={expected_workers})" ) break time.sleep(0.1) else: raise AssertionError( f"compaction did not expose chunks={expected_chunks}, " f"workers={expected_workers}\n{log}" )