#!/usr/bin/env python3 """Offline process receipt for Jetstream V2 log level and format semantics.""" from __future__ import annotations import json import signal import socket import subprocess import tempfile import time import urllib.request BIN = "./zig-out/bin/stream" def unused_port() -> int: sock = socket.socket() sock.bind(("127.0.0.1", 0)) port = sock.getsockname()[1] sock.close() return port def start(root: str, port: int, prefix: list[str]) -> subprocess.Popen[str]: return subprocess.Popen( [ BIN, *prefix, f"--addr=127.0.0.1:{port}", f"--data-dir={root}", "--upstream=ws://127.0.0.1:17990", "--relay-http=http://127.0.0.1:17990", "--plc-url=http://127.0.0.1:17990", "--compaction-interval=0", "--retry-interval=0", "--no-verify", ], stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, text=True, ) def wait_ready(process: subprocess.Popen[str], port: int) -> None: # listener-first startup: the listener answers from bind time, so wait # for /status to leave "starting" (storage published, pipeline starting) deadline = time.monotonic() + 10 while time.monotonic() < deadline: if process.poll() is not None: stderr = process.stderr.read() if process.stderr else "" raise AssertionError(f"stream exited with {process.returncode}:\n{stderr}") try: body = urllib.request.urlopen(f"http://127.0.0.1:{port}/status", timeout=0.2).read() if b"state starting" not in body: return time.sleep(0.025) except OSError: time.sleep(0.025) raise AssertionError(f"listener {port} did not start") def stop(process: subprocess.Popen[str]) -> list[str]: process.send_signal(signal.SIGTERM) process.wait(timeout=5) assert process.returncode == 0, process.stderr.read() if process.stderr else "" return [line for line in process.stderr.read().splitlines() if line] def default_json_case(root: str) -> None: port = unused_port() process = start(root, port, ["serve"]) wait_ready(process, port) records = [json.loads(line) for line in stop(process)] assert records assert all(set(("time", "level", "msg", "component")) <= record.keys() for record in records) assert any(record["level"] == "INFO" for record in records) assert not any(record["level"] == "DEBUG" for record in records) def warn_json_case(root: str) -> None: port = unused_port() process = start(root, port, ["serve", "--log-level=warning"]) wait_ready(process, port) # the case needs at least one failed upstream dial (the ERROR record) # before shutdown; with listener-first startup the listener answers # before the pipeline's first dial, so wait for the record itself # instead of sleeping a fixed slice import threading lines: list[str] = [] got_error = threading.Event() def pump() -> None: assert process.stderr is not None for line in process.stderr: if line.strip(): lines.append(line) try: if json.loads(line)["level"] == "ERROR": got_error.set() except ValueError: pass reader = threading.Thread(target=pump, daemon=True) reader.start() assert got_error.wait(10), f"no ERROR record before deadline: {lines!r}" process.send_signal(signal.SIGTERM) process.wait(timeout=5) assert process.returncode == 0 reader.join(timeout=5) records = [json.loads(line) for line in lines] assert records assert any(record["level"] == "ERROR" for record in records) assert all(record["level"] in ("WARN", "ERROR") for record in records) def root_flags_text_debug_case(root: str) -> None: port = unused_port() process = start( root, port, ["--log-level=DBG", "--log-format=text", "serve"], ) wait_ready(process, port) lines = stop(process) assert lines assert all(line.startswith("time=") for line in lines) assert any(" level=DEBUG " in line for line in lines) assert any('msg="logging configured: level=DBG format=text"' in line for line in lines) def invalid_config_case() -> None: for flag, expected in ( ("--log-level=loud", "InvalidLogLevel"), ("--log-format=xml", "InvalidLogFormat"), ): result = subprocess.run( [BIN, "serve", flag], stdout=subprocess.DEVNULL, stderr=subprocess.PIPE, text=True, timeout=3, check=False, ) assert result.returncode != 0 assert expected in result.stderr, result.stderr def persistent_command_flags_case() -> None: for args in ( ["--log-level", "debug", "--log-format", "text", "version"], ["version", "--log-level=debug", "--log-format=text"], ): result = subprocess.run( [BIN, *args], stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, timeout=3, check=False, ) assert result.returncode == 0, result.stderr assert result.stdout.startswith("jetstream version "), result.stdout def main() -> None: with tempfile.TemporaryDirectory() as root: default_json_case(f"{root}/default") warn_json_case(f"{root}/warn") root_flags_text_debug_case(f"{root}/debug") persistent_command_flags_case() invalid_config_case() print("logging contract: json-default warn-filter text-debug persistent-flags invalid-config passed") if __name__ == "__main__": main()